Skip to content

Commit 15a7ce1

Browse files
authored
IcebergIO and AddFiles multi bucket catalog validation (#39907)
* Validate multi region biglake catalog * trigger tests * remove x-region test * fix tests, cleanup test artifacts * fixes * rename
1 parent 4203bab commit 15a7ce1

8 files changed

Lines changed: 491 additions & 119 deletions

File tree

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
11
{
22
"comment": "Modify this file in a trivial way to cause this test suite to run.",
3-
"modification": 10
3+
"modification": 11
44
}

‎sdks/java/io/iceberg/build.gradle‎

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -179,10 +179,15 @@ task integrationTest(type: Test) {
179179
"--project=${gcpProject}",
180180
"--tempLocation=${gcpTempLocation}",
181181
])
182-
// Warehouse (= catalog) used by the BigLake REST catalog tests; overridable for runs
183-
// against a non-default project's catalog.
184-
systemProperty "beam.iceberg.biglake.warehouse",
185-
project.findProperty('biglakeWarehouse') ?: 'gs://managed-iceberg-biglake-its'
182+
// Multiple-bucket Lakehouse REST catalog used by RESTCatalogBLMSIT and AddFilesIT (see
183+
// LakehouseTestCatalog). Locations: the catalog's default location first, then a restricted
184+
// location in another bucket. Override for runs against another project's catalog:
185+
// -PlakehouseWarehouse=bl://projects/PROJECT/catalogs/CATALOG
186+
// -PlakehouseLocations=gs://default-bucket/path,gs://other-bucket/path
187+
systemProperty "beam.iceberg.lakehouse.warehouse",
188+
project.findProperty('lakehouseWarehouse') ?: 'bl://projects/apache-beam-testing/catalogs/beam-lakehouse-it'
189+
systemProperty "beam.iceberg.lakehouse.locations",
190+
project.findProperty('lakehouseLocations') ?: 'gs://beam-lakehouse-it,gs://beam-lakehouse-it-added-path'
186191
// Connection + storage root for BigQueryManagedTableCrossEngineIT.
187192
if (project.findProperty('bqImtConnection') != null) {
188193
systemProperty "beam.bq.imt.connection", project.findProperty('bqImtConnection')

‎sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesIT.java‎

Lines changed: 165 additions & 54 deletions
Large diffs are not rendered by default.
Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
package org.apache.beam.sdk.io.iceberg;
19+
20+
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
21+
22+
import com.google.api.services.storage.model.StorageObject;
23+
import java.io.IOException;
24+
import java.util.ArrayList;
25+
import java.util.List;
26+
import java.util.Map;
27+
import java.util.regex.Matcher;
28+
import java.util.regex.Pattern;
29+
import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
30+
import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
31+
import org.apache.beam.sdk.extensions.gcp.util.GcsUtil;
32+
import org.apache.beam.sdk.testing.TestPipeline;
33+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter;
34+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
35+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
36+
import org.apache.iceberg.catalog.Catalog;
37+
import org.apache.iceberg.catalog.Namespace;
38+
import org.apache.iceberg.catalog.SupportsNamespaces;
39+
import org.apache.iceberg.catalog.TableIdentifier;
40+
41+
/**
42+
* Test-side description of the Lakehouse Iceberg REST catalog the ITs run against.
43+
*
44+
* <p>The catalog is a multiple-bucket catalog addressed by a {@code
45+
* bl://projects/PROJECT/catalogs/CATALOG} warehouse and allowed to store resources under a fixed
46+
* set of Cloud Storage locations: the catalog's default location plus any restricted locations.
47+
* Configured (defaults and overrides) by the integrationTest task in build.gradle through the
48+
* system properties
49+
*
50+
* <ul>
51+
* <li>{@code beam.iceberg.lakehouse.warehouse}: the {@code bl://} warehouse URI
52+
* <li>{@code beam.iceberg.lakehouse.locations}: comma-separated {@code gs://} prefixes the
53+
* catalog may write to; the first is the catalog's default location, the rest are additional
54+
* restricted locations
55+
* </ul>
56+
*/
57+
public final class LakehouseTestCatalog {
58+
private static final Pattern WAREHOUSE_PATTERN =
59+
Pattern.compile("bl://projects/([^/]+)/catalogs/([^/]+)");
60+
61+
public static final String WAREHOUSE = requiredProperty("beam.iceberg.lakehouse.warehouse");
62+
63+
public static final List<String> LOCATIONS =
64+
ImmutableList.copyOf(
65+
Splitter.on(',')
66+
.trimResults()
67+
.omitEmptyStrings()
68+
.split(requiredProperty("beam.iceberg.lakehouse.locations")));
69+
70+
/** Catalog id, which is also the second segment of BigQuery's 4-part table reference. */
71+
public static final String CATALOG_ID = parseCatalogId(WAREHOUSE);
72+
73+
private static final String PROJECT =
74+
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
75+
76+
private LakehouseTestCatalog() {}
77+
78+
private static String requiredProperty(String name) {
79+
String value = System.getProperty(name);
80+
checkArgument(
81+
value != null && !value.isEmpty(),
82+
"System property %s is not set; run through the integrationTest gradle task, which"
83+
+ " sets it, or pass -D%s=...",
84+
name,
85+
name);
86+
return value;
87+
}
88+
89+
private static String parseCatalogId(String warehouse) {
90+
// Legacy single-bucket catalogs are addressed by their bucket and named after it.
91+
if (warehouse.startsWith("gs://")) {
92+
return bucketOf(warehouse);
93+
}
94+
Matcher matcher = WAREHOUSE_PATTERN.matcher(warehouse);
95+
checkArgument(
96+
matcher.matches(),
97+
"Expected a bl://projects/PROJECT/catalogs/CATALOG (or legacy gs://BUCKET) warehouse, got '%s'",
98+
warehouse);
99+
return matcher.group(2);
100+
}
101+
102+
/** The catalog's default location; tables land here unless created with an explicit location. */
103+
public static String defaultLocation() {
104+
return LOCATIONS.get(0);
105+
}
106+
107+
/** A restricted location outside the default one, i.e. in a second bucket. */
108+
public static String additionalLocation() {
109+
checkArgument(
110+
LOCATIONS.size() >= 2,
111+
"beam.iceberg.lakehouse.locations must list at least two locations, got %s",
112+
LOCATIONS);
113+
return LOCATIONS.get(1);
114+
}
115+
116+
public static String bucketOf(String gcsLocation) {
117+
checkArgument(gcsLocation.startsWith("gs://"), "Not a gs:// location: %s", gcsLocation);
118+
String withoutScheme = gcsLocation.substring("gs://".length());
119+
int slash = withoutScheme.indexOf('/');
120+
return slash < 0 ? withoutScheme : withoutScheme.substring(0, slash);
121+
}
122+
123+
/** Object-name prefix (no bucket, no leading slash) of a {@code gs://bucket/path} location. */
124+
public static String prefixOf(String gcsLocation) {
125+
String withoutScheme = gcsLocation.substring("gs://".length());
126+
int slash = withoutScheme.indexOf('/');
127+
return slash < 0 ? "" : withoutScheme.substring(slash + 1);
128+
}
129+
130+
public static Map<String, String> catalogProperties() {
131+
return ImmutableMap.<String, String>builder()
132+
.put("type", "rest")
133+
.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog")
134+
.put("warehouse", WAREHOUSE)
135+
.put("header.x-goog-user-project", PROJECT)
136+
// Required by catalogs in vended-credentials mode; ignored in end-user mode.
137+
.put("header.X-Iceberg-Access-Delegation", "vended-credentials")
138+
.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO")
139+
.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager")
140+
.build();
141+
}
142+
143+
/**
144+
* Drops every table in the namespaces and then the namespaces themselves, and deletes the tables'
145+
* files from Cloud Storage.
146+
*/
147+
public static void dropNamespacesAndFiles(Catalog catalog, List<String> namespaces)
148+
throws IOException {
149+
for (String name : namespaces) {
150+
Namespace namespace = Namespace.of(name);
151+
if (!((SupportsNamespaces) catalog).namespaceExists(namespace)) {
152+
continue;
153+
}
154+
dropTablesAndFiles(catalog, name);
155+
((SupportsNamespaces) catalog).dropNamespace(namespace);
156+
}
157+
}
158+
159+
/**
160+
* Drops every table in the namespace and deletes the tables' files from Cloud Storage: Lakehouse
161+
* keeps a dropped table's data and metadata (even with purge), and table locations carry a random
162+
* suffix, so they are captured before the drop.
163+
*/
164+
public static void dropTablesAndFiles(Catalog catalog, String namespaceName) throws IOException {
165+
Namespace namespace = Namespace.of(namespaceName);
166+
if (!((SupportsNamespaces) catalog).namespaceExists(namespace)) {
167+
return;
168+
}
169+
List<String> tableLocations = new ArrayList<>();
170+
for (TableIdentifier identifier : catalog.listTables(namespace)) {
171+
tableLocations.add(catalog.loadTable(identifier).location());
172+
catalog.dropTable(identifier);
173+
}
174+
for (String location : tableLocations) {
175+
deleteObjects(location);
176+
}
177+
}
178+
179+
/** Deletes every object under a {@code gs://bucket/prefix} location. */
180+
public static void deleteObjects(String gcsLocation) throws IOException {
181+
GcsUtil gcsUtil = TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
182+
List<StorageObject> objects =
183+
gcsUtil.listObjects(bucketOf(gcsLocation), prefixOf(gcsLocation), null).getItems();
184+
if (objects == null || objects.isEmpty()) {
185+
return;
186+
}
187+
List<String> paths = new ArrayList<>();
188+
for (StorageObject object : objects) {
189+
paths.add("gs://" + object.getBucket() + "/" + object.getName());
190+
}
191+
gcsUtil.remove(paths);
192+
}
193+
194+
/** BigQuery's 4-part {@code project.catalog.namespace.table} reference. */
195+
public static String bigQueryTableSpec(String namespace, String table) {
196+
return String.format("%s.%s.%s.%s", PROJECT, CATALOG_ID, namespace, table);
197+
}
198+
}

‎sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCatalogBaseIT.java‎

Lines changed: 10 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -175,10 +175,9 @@ public abstract class IcebergCatalogBaseIT implements Serializable {
175175
/**
176176
* Catalogs whose tables are also queryable with BigQuery return the BigQuery table reference for
177177
* the given Iceberg table id: either the 4-part {@code project.catalog.namespace.table} form for
178-
* Lakehouse runtime catalog (BigLake metastore REST) tables, or the 3-part {@code
179-
* project.dataset.table} form for the BigQuery metastore federation, where namespaces surface as
180-
* datasets. Returning null (the default) disables the cross-engine read checks in {@link
181-
* #testReadWithBigQueryIO()}.
178+
* Lakehouse runtime catalog (Iceberg REST) tables, or the 3-part {@code project.dataset.table}
179+
* form for the BigQuery metastore federation, where namespaces surface as datasets. Returning
180+
* null (the default) disables the cross-engine read checks in {@link #testReadWithBigQueryIO()}.
182181
*/
183182
public @Nullable String bigQueryTableSpec(String tableId) {
184183
return null;
@@ -238,14 +237,14 @@ public void cleanUp() throws Exception {
238237
try {
239238
GcsUtil gcsUtil = OPTIONS.as(GcsOptions.class).getGcsUtil();
240239
GcsPath path = GcsPath.fromUri(warehouse);
240+
// The warehouse may be a bare bucket (no object path), where getFileName() throws.
241+
String prefix =
242+
path.getObject().isEmpty()
243+
? getClass().getSimpleName()
244+
: getClass().getSimpleName() + "/" + path.getFileName();
241245

242246
@Nullable List<StorageObject> objects =
243-
gcsUtil
244-
.listObjects(
245-
path.getBucket(),
246-
getClass().getSimpleName() + "/" + path.getFileName().toString(),
247-
null)
248-
.getItems();
247+
gcsUtil.listObjects(path.getBucket(), prefix, null).getItems();
249248

250249
// sometimes a catalog's cleanup will take care of all the files.
251250
// If any files are left though, manually delete them with GCS utils
@@ -432,7 +431,7 @@ private List<Row> populateTable(Table table, @Nullable String charOverride) thro
432431
}
433432
}
434433

435-
private List<Record> readRecords(Table table) throws IOException {
434+
protected List<Record> readRecords(Table table) throws IOException {
436435
org.apache.iceberg.Schema tableSchema = table.schema();
437436
TableScan tableScan = table.newScan().project(tableSchema);
438437
List<Record> writtenRecords = new ArrayList<>();

‎sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCdcWriteBaseIT.java‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -236,13 +236,13 @@ public void cleanUp() throws Exception {
236236
try {
237237
GcsUtil gcsUtil = OPTIONS.as(GcsOptions.class).getGcsUtil();
238238
GcsPath path = GcsPath.fromUri(warehouse);
239+
// The warehouse may be a bare bucket (no object path), where getFileName() throws.
240+
String prefix =
241+
path.getObject().isEmpty()
242+
? getClass().getSimpleName()
243+
: getClass().getSimpleName() + "/" + path.getFileName();
239244
@Nullable List<StorageObject> objects =
240-
gcsUtil
241-
.listObjects(
242-
path.getBucket(),
243-
getClass().getSimpleName() + "/" + path.getFileName().toString(),
244-
null)
245-
.getItems();
245+
gcsUtil.listObjects(path.getBucket(), prefix, null).getItems();
246246
// A catalog's cleanup sometimes removes every file; delete whatever is left.
247247
if (objects != null) {
248248
gcsUtil.remove(

‎sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/LakehouseCatalogCdcWriteIT.java‎

Lines changed: 9 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,9 @@
1717
*/
1818
package org.apache.beam.sdk.io.iceberg.catalog;
1919

20+
import java.io.IOException;
2021
import java.util.Map;
22+
import org.apache.beam.sdk.io.iceberg.LakehouseTestCatalog;
2123
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
2224
import org.apache.iceberg.catalog.Catalog;
2325
import org.apache.iceberg.rest.RESTCatalog;
@@ -28,27 +30,19 @@
2830
public class LakehouseCatalogCdcWriteIT extends IcebergCdcWriteBaseIT {
2931
private static Map<String, String> catalogProps;
3032

31-
private static final String LAKEHOUSE_WAREHOUSE =
32-
System.getProperty("beam.iceberg.biglake.warehouse", "gs://managed-iceberg-biglake-its");
33-
3433
@BeforeClass
3534
public static void setup() {
36-
warehouse = LAKEHOUSE_WAREHOUSE;
37-
catalogProps =
38-
ImmutableMap.<String, String>builder()
39-
.put("type", "rest")
40-
.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog")
41-
.put("warehouse", LAKEHOUSE_WAREHOUSE)
42-
.put("header.x-goog-user-project", OPTIONS.getProject())
43-
.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO")
44-
.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager")
45-
.build();
35+
warehouse = LakehouseTestCatalog.defaultLocation();
36+
catalogProps = LakehouseTestCatalog.catalogProperties();
4637
}
4738

4839
@After
49-
public void after() {
40+
public void after() throws IOException {
41+
// Lakehouse keeps a dropped table's files, so remove them before the base class drops the
42+
// namespace.
43+
LakehouseTestCatalog.dropTablesAndFiles(catalog, namespace());
5044
// The base class points its cleanup at this warehouse.
51-
warehouse = LAKEHOUSE_WAREHOUSE;
45+
warehouse = LakehouseTestCatalog.defaultLocation();
5246
}
5347

5448
@Override

0 commit comments

Comments
 (0)