This is an automated email from the ASF dual-hosted git repository.
claudevdm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 15a7ce17e1f IcebergIO and AddFiles multi bucket catalog validation
(#39907)
15a7ce17e1f is described below
commit 15a7ce17e1fb22897ac24527ad7eaa37ab75f786
Author: claudevdm <[email protected]>
AuthorDate: Fri Oct 2 10:33:16 2026 -0400
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
---
.../IO_Iceberg_Integration_Tests.json | 2 +-
sdks/java/io/iceberg/build.gradle | 13 +-
.../org/apache/beam/sdk/io/iceberg/AddFilesIT.java | 219 ++++++++++++++++-----
.../beam/sdk/io/iceberg/LakehouseTestCatalog.java | 198 +++++++++++++++++++
.../io/iceberg/catalog/IcebergCatalogBaseIT.java | 21 +-
.../io/iceberg/catalog/IcebergCdcWriteBaseIT.java | 12 +-
.../catalog/LakehouseCatalogCdcWriteIT.java | 24 +--
.../sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java | 121 +++++++++---
8 files changed, 491 insertions(+), 119 deletions(-)
diff --git a/.github/trigger_files/IO_Iceberg_Integration_Tests.json
b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
index 7392be3b11c..e1290fa50c3 100644
--- a/.github/trigger_files/IO_Iceberg_Integration_Tests.json
+++ b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
@@ -1,4 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to
run.",
- "modification": 10
+ "modification": 11
}
diff --git a/sdks/java/io/iceberg/build.gradle
b/sdks/java/io/iceberg/build.gradle
index 596962b1a03..86964472cd8 100644
--- a/sdks/java/io/iceberg/build.gradle
+++ b/sdks/java/io/iceberg/build.gradle
@@ -179,10 +179,15 @@ task integrationTest(type: Test) {
"--project=${gcpProject}",
"--tempLocation=${gcpTempLocation}",
])
- // Warehouse (= catalog) used by the BigLake REST catalog tests;
overridable for runs
- // against a non-default project's catalog.
- systemProperty "beam.iceberg.biglake.warehouse",
- project.findProperty('biglakeWarehouse') ?:
'gs://managed-iceberg-biglake-its'
+ // Multiple-bucket Lakehouse REST catalog used by RESTCatalogBLMSIT and
AddFilesIT (see
+ // LakehouseTestCatalog). Locations: the catalog's default location first,
then a restricted
+ // location in another bucket. Override for runs against another project's
catalog:
+ // -PlakehouseWarehouse=bl://projects/PROJECT/catalogs/CATALOG
+ // -PlakehouseLocations=gs://default-bucket/path,gs://other-bucket/path
+ systemProperty "beam.iceberg.lakehouse.warehouse",
+ project.findProperty('lakehouseWarehouse') ?:
'bl://projects/apache-beam-testing/catalogs/beam-lakehouse-it'
+ systemProperty "beam.iceberg.lakehouse.locations",
+ project.findProperty('lakehouseLocations') ?:
'gs://beam-lakehouse-it,gs://beam-lakehouse-it-added-path'
// Connection + storage root for BigQueryManagedTableCrossEngineIT.
if (project.findProperty('bqImtConnection') != null) {
systemProperty "beam.bq.imt.connection",
project.findProperty('bqImtConnection')
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesIT.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesIT.java
index 5afc6107c58..a1bb7943175 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesIT.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesIT.java
@@ -22,11 +22,15 @@ import static
org.apache.beam.sdk.io.FileIO.Write.defaultNaming;
import static
org.apache.beam.sdk.io.iceberg.IcebergUtils.beamSchemaToIcebergSchema;
import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
import static org.apache.beam.sdk.values.TypeDescriptors.strings;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.startsWith;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
+import com.google.api.client.http.HttpResponseException;
import com.google.api.services.storage.model.StorageObject;
import com.google.cloud.storage.Blob;
import com.google.cloud.storage.Notification;
@@ -79,6 +83,7 @@ import org.apache.beam.sdk.values.TypeDescriptor;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
import org.apache.hadoop.util.Lists;
+import org.apache.iceberg.BaseTable;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Snapshot;
@@ -105,12 +110,11 @@ import org.slf4j.LoggerFactory;
public class AddFilesIT {
private static final Logger LOG = LoggerFactory.getLogger(AddFilesIT.class);
- // Bucket-backed BigLake catalogs are named after their bucket. Overridable
for local runs
- // against a different project's catalog:
-Dbeam.iceberg.biglake.warehouse=gs://my-bucket
- private static final String CATALOG_NAME =
- System.getProperty("beam.iceberg.biglake.warehouse",
"gs://managed-iceberg-biglake-its")
- .replace("gs://", "");
- private static final String WAREHOUSE = "gs://" + CATALOG_NAME;
+ // Multiple-bucket Lakehouse catalog (see LakehouseTestCatalog). Source
parquet files, and the
+ // GCS notifications announcing them, live under the catalog's default
location.
+ private static final String DATA_LOCATION =
LakehouseTestCatalog.defaultLocation();
+ private static final String DATA_BUCKET =
LakehouseTestCatalog.bucketOf(DATA_LOCATION);
+ private static final String DATA_PREFIX =
LakehouseTestCatalog.prefixOf(DATA_LOCATION);
private static final String PROJECT =
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
@Rule public TestName testName = new TestName();
@@ -127,20 +131,14 @@ public class AddFilesIT {
.addStringField("name")
.addStringField("kind")
.build();
- private static final Map<String, String> BIGLAKE_PROPS =
- Map.of(
- "type", "rest",
- "uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog",
- "warehouse", WAREHOUSE,
- "header.x-goog-user-project", PROJECT,
- "rest.auth.type", "google",
- "io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO",
- // required by vended-credentials catalogs, harmless on end-user ones
- "header.X-Iceberg-Access-Delegation", "vended-credentials");
+ private static final Map<String, String> LAKEHOUSE_PROPS =
+ LakehouseTestCatalog.catalogProperties();
private Storage storage;
private PubsubClient pubsub;
private Notification notification;
private final String namespace = getClass().getSimpleName() + "_" +
System.currentTimeMillis();
+ // Namespace placed in the catalog's additional location (second bucket).
+ private final String altNamespace = namespace + "_alt";
private String srcTableName;
private String destTableName;
private TableIdentifier srcTableId;
@@ -176,45 +174,77 @@ public class AddFilesIT {
.setPayloadFormat(NotificationInfo.PayloadFormat.JSON_API_V1)
.build();
try {
- notification = storage.createNotification(WAREHOUSE.replace("gs://",
""), notificationInfo);
+ notification = storage.createNotification(DATA_BUCKET, notificationInfo);
} catch (StorageException e) {
if (e.getMessage().contains("Too many overlapping notifications")) {
- List<Notification> existing =
storage.listNotifications(WAREHOUSE.replace("gs://", ""));
+ List<Notification> existing = storage.listNotifications(DATA_BUCKET);
LOG.warn(
"Too many notifications on bucket {}: {}. Deleting existing
notifications to make room: {}",
- WAREHOUSE,
+ DATA_BUCKET,
e,
existing.stream()
.map(NotificationInfo::getNotificationId)
.collect(Collectors.toList()));
- existing.forEach(
- n -> storage.deleteNotification(WAREHOUSE.replace("gs://", ""),
n.getNotificationId()));
+ existing.forEach(n -> storage.deleteNotification(DATA_BUCKET,
n.getNotificationId()));
// try creating it again
- notification = storage.createNotification(WAREHOUSE.replace("gs://",
""), notificationInfo);
+ notification = storage.createNotification(DATA_BUCKET,
notificationInfo);
} else {
+ logNotificationFailure(e);
throw e;
}
}
salt = System.currentTimeMillis();
- dirName = format("%s-%s/%s", getClass().getSimpleName(), salt,
testName.getMethodName());
+ // Object-name prefix of this test's parquet files; DATA_PREFIX is empty
for a bare bucket.
+ dirName =
+ format(
+ "%s%s-%s/%s",
+ DATA_PREFIX.isEmpty() ? "" : DATA_PREFIX + "/",
+ getClass().getSimpleName(),
+ salt,
+ testName.getMethodName());
srcTableName = "src_" + testName.getMethodName() + "_" + salt;
destTableName = "dest_" + testName.getMethodName() + "_" + salt;
srcTableId = TableIdentifier.of(namespace, srcTableName);
destTableId = TableIdentifier.of(namespace, destTableName);
- catalog.initialize("test_catalog", BIGLAKE_PROPS);
+ catalog.initialize("test_catalog", LAKEHOUSE_PROPS);
cleanupCatalog();
catalog.createNamespace(Namespace.of(namespace));
}
- private void cleanupCatalog() {
- Namespace ns = Namespace.of(namespace);
- if (catalog.namespaceExists(ns)) {
- catalog.listTables(ns).forEach(catalog::dropTable);
- catalog.dropNamespace(ns);
+ /**
+ * GCS reports server-side errors (5xx) without a cause, so log what
identifies the request to GCS
+ * support and the bucket's notification count at the time.
+ */
+ private void logNotificationFailure(StorageException e) {
+ @Nullable String requestId = null;
+ Throwable cause = e.getCause();
+ if (cause instanceof HttpResponseException) {
+ HttpResponseException response = (HttpResponseException) cause;
+ requestId =
response.getHeaders().getFirstHeaderStringValue("x-guploader-uploadid");
+ }
+ String existingNotifications;
+ try {
+ existingNotifications =
String.valueOf(storage.listNotifications(DATA_BUCKET).size());
+ } catch (StorageException listError) {
+ existingNotifications = "unknown (" + listError.getMessage() + ")";
}
+ LOG.error(
+ "Failed to create a GCS notification on bucket {} for topic {}: HTTP
{}, reason={},"
+ + " retryable={}, x-guploader-uploadid={}, notifications on
bucket={}",
+ DATA_BUCKET,
+ notificationsTopic,
+ e.getCode(),
+ e.getReason(),
+ e.isRetryable(),
+ requestId,
+ existingNotifications);
+ }
+
+ private void cleanupCatalog() throws IOException {
+ LakehouseTestCatalog.dropNamespacesAndFiles(catalog,
Arrays.asList(namespace, altNamespace));
}
@After
@@ -227,7 +257,7 @@ public class AddFilesIT {
}
try {
- storage.deleteNotification(WAREHOUSE.replace("gs://", ""),
notification.getNotificationId());
+ storage.deleteNotification(DATA_BUCKET,
notification.getNotificationId());
storage.close();
} catch (Exception e) {
LOG.warn("Failed to clean up GCS notifications", e);
@@ -241,9 +271,7 @@ public class AddFilesIT {
try {
Iterable<Blob> blobs =
- storage
- .list(WAREHOUSE.replace("gs://", ""),
Storage.BlobListOption.prefix(dirName))
- .getValues();
+ storage.list(DATA_BUCKET,
Storage.BlobListOption.prefix(dirName)).getValues();
blobs.forEach(b -> storage.delete(b.getBlobId()));
} catch (Exception e) {
LOG.warn("Failed to clean up GCS bucket", e);
@@ -256,7 +284,7 @@ public class AddFilesIT {
// first create a source iceberg table
catalog.createTable(srcTableId, beamSchemaToIcebergSchema(ROW_SCHEMA),
SPEC);
- // BigLake may write under {namespace}/{table}/{id}/data/... rather than
the Hive-style
+ // Lakehouse may write under {namespace}/{table}/{id}/data/... rather than
the Hive-style
// {namespace}/{table}/data/... layout, so match the table prefix and a
/data/ segment.
String tablePrefix = format("%s/%s/", namespace, srcTableName);
@@ -276,7 +304,7 @@ public class AddFilesIT {
Managed.write(Managed.ICEBERG)
.withConfig(
ImmutableMap.of(
- "table", srcTableId.toString(), "catalog_properties",
BIGLAKE_PROPS)));
+ "table", srcTableId.toString(), "catalog_properties",
LAKEHOUSE_PROPS)));
q.run().waitUntilFinish();
// check that the destination table has been created
@@ -352,8 +380,8 @@ public class AddFilesIT {
throws InterruptedException, TimeoutException, IOException {
// start with a table that does not exist
- String parquetDir = format("%s/%s/", WAREHOUSE, dirName);
- String tempDir = format("%s/%s-tmp/", WAREHOUSE, dirName);
+ String parquetDir = format("gs://%s/%s/", DATA_BUCKET, dirName);
+ String tempDir = format("gs://%s/%s-tmp/", DATA_BUCKET, dirName);
// let the add files pipeline run in the background
PipelineResult addFilesPipeline = startAddFilesListener(dirName);
@@ -385,8 +413,7 @@ public class AddFilesIT {
GcsUtil gcsUtil =
TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
- Iterable<StorageObject> objects =
- gcsUtil.listObjects(WAREHOUSE.replace("gs://", ""), dirName,
null).getItems();
+ Iterable<StorageObject> objects = gcsUtil.listObjects(DATA_BUCKET,
dirName, null).getItems();
List<String> writtenFilePaths =
Lists.newArrayList(objects).stream()
.map(o -> format("gs://%s/%s", o.getBucket(), o.getName()))
@@ -426,11 +453,92 @@ public class AddFilesIT {
testBatchParquetImport(true);
}
+ /**
+ * The destination table lives in the catalog's additional location (a
second bucket) while the
+ * source parquet files stay in the default one. Lakehouse pins tables under
their namespace's
+ * location, so the table is created in a namespace placed in the second
bucket; AddFiles must
+ * commit metadata there and reference the files in place across buckets.
+ */
+ @Test
+ public void testBatchParquetImportToTableInAdditionalLocation() throws
IOException {
+ String namespaceLocation = LakehouseTestCatalog.additionalLocation() + "/"
+ altNamespace;
+ assertNotEquals(
+ "Test needs two distinct buckets",
+ DATA_BUCKET,
+ LakehouseTestCatalog.bucketOf(namespaceLocation));
+ catalog.createNamespace(
+ Namespace.of(altNamespace), ImmutableMap.of("location",
namespaceLocation));
+ destTableId = TableIdentifier.of(altNamespace, destTableName);
+ catalog.createTable(destTableId, beamSchemaToIcebergSchema(ROW_SCHEMA),
SPEC);
+ assertThat(catalog.loadTable(destTableId).location(),
startsWith(namespaceLocation));
+
+ List<String> writtenFilePaths = writeParquetFiles();
+ Pipeline p = Pipeline.create();
+ PCollectionRowTuple tuple =
+ p.apply(Create.of(writtenFilePaths))
+ .apply(
+ new AddFiles(
+
IcebergCatalogConfig.builder().setCatalogProperties(LAKEHOUSE_PROPS).build(),
+ destTableId.toString(),
+ null,
+ PARTITION_FIELDS,
+ null,
+ TABLE_PROPS,
+ null,
+ null));
+ PAssert.that(tuple.get("errors")).empty();
+ p.run().waitUntilFinish();
+
+ assertTrue(checkTableHasRegisteredParquetFiles(writtenFilePaths));
+ Table destTable = catalog.loadTable(destTableId);
+ String metadataLocation = ((BaseTable)
destTable).operations().current().metadataFileLocation();
+ assertThat(metadataLocation,
startsWith(LakehouseTestCatalog.additionalLocation()));
+ for (String path : writtenFilePaths) {
+ assertThat(path, startsWith("gs://" + DATA_BUCKET + "/"));
+ }
+ checkRecordsInDestinationTable(/* alsoCheckWithBigQueryIO= */ true);
+ }
+
+ /** Writes TEST_ROWS as parquet under the test's data dir and returns the
written file paths. */
+ private List<String> writeParquetFiles() throws IOException {
+ String parquetDir = format("gs://%s/%s/", DATA_BUCKET, dirName);
+ String tempDir = format("gs://%s/%s-tmp/", DATA_BUCKET, dirName);
+ LOG.info("Writing records to the parquet dir");
+ Pipeline q = Pipeline.create();
+ org.apache.avro.Schema avroSchema = AvroUtils.toAvroSchema(ROW_SCHEMA);
+ q.apply(Create.of(TEST_ROWS))
+ .setRowSchema(ROW_SCHEMA)
+ .apply(
+ MapElements.into(TypeDescriptor.of(GenericRecord.class))
+ .via(AvroUtils.getRowToGenericRecordFunction(avroSchema)))
+ .setCoder(AvroCoder.of(avroSchema))
+ .apply(
+ FileIO.<String, GenericRecord>writeDynamic()
+ .by(
+ record ->
+ format("%s-%s-%s", record.get("id"),
record.get("name"), record.get("age")))
+ .via(ParquetIO.sink(avroSchema))
+ .withNaming(name -> defaultNaming(name, ".parquet"))
+ .withTempDirectory(tempDir)
+ .to(parquetDir)
+ .withDestinationCoder(StringUtf8Coder.of()));
+ q.run().waitUntilFinish();
+
+ GcsUtil gcsUtil =
TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
+ Iterable<StorageObject> objects = gcsUtil.listObjects(DATA_BUCKET,
dirName, null).getItems();
+ List<String> writtenFilePaths =
+ Lists.newArrayList(objects).stream()
+ .map(o -> format("gs://%s/%s", o.getBucket(), o.getName()))
+ .collect(Collectors.toList());
+ LOG.info("Written file paths: {}", writtenFilePaths);
+ return writtenFilePaths;
+ }
+
private void testBatchParquetImport(boolean isUIT) throws IOException {
// start with a table that does not exist
- String parquetDir = format("%s/%s/", WAREHOUSE, dirName);
- String tempDir = format("%s/%s-tmp/", WAREHOUSE, dirName);
+ String parquetDir = format("gs://%s/%s/", DATA_BUCKET, dirName);
+ String tempDir = format("gs://%s/%s-tmp/", DATA_BUCKET, dirName);
// write some parquet files
LOG.info("Writing records to the parquet dir");
@@ -456,8 +564,7 @@ public class AddFilesIT {
GcsUtil gcsUtil =
TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
- Iterable<StorageObject> objects =
- gcsUtil.listObjects(WAREHOUSE.replace("gs://", ""), dirName,
null).getItems();
+ Iterable<StorageObject> objects = gcsUtil.listObjects(DATA_BUCKET,
dirName, null).getItems();
List<String> writtenFilePaths =
Lists.newArrayList(objects).stream()
.map(o -> format("gs://%s/%s", o.getBucket(), o.getName()))
@@ -478,7 +585,7 @@ public class AddFilesIT {
p.apply(Create.of(writtenFilePaths))
.apply(
new AddFiles(
-
IcebergCatalogConfig.builder().setCatalogProperties(BIGLAKE_PROPS).build(),
+
IcebergCatalogConfig.builder().setCatalogProperties(LAKEHOUSE_PROPS).build(),
namespace + "." + destTableName,
null,
isUIT ? null : PARTITION_FIELDS,
@@ -511,15 +618,13 @@ public class AddFilesIT {
Schema narrow =
Schema.builder().addInt64Field("id").addStringField("name").build();
catalog.createTable(destTableId, beamSchemaToIcebergSchema(narrow));
- String parquetDir = format("%s/%s/", WAREHOUSE, dirName);
- String tempDir = format("%s/%s-tmp/", WAREHOUSE, dirName);
+ String parquetDir = format("gs://%s/%s/", DATA_BUCKET, dirName);
+ String tempDir = format("gs://%s/%s-tmp/", DATA_BUCKET, dirName);
writeParquet(TEST_ROWS, ROW_SCHEMA, parquetDir + "plain/", tempDir +
"plain/");
GcsUtil gcsUtil =
TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
List<String> writtenFilePaths =
- Lists.newArrayList(
- gcsUtil.listObjects(WAREHOUSE.replace("gs://", ""), dirName,
null).getItems())
- .stream()
+ Lists.newArrayList(gcsUtil.listObjects(DATA_BUCKET, dirName,
null).getItems()).stream()
.map(o -> format("gs://%s/%s", o.getBucket(), o.getName()))
.collect(Collectors.toList());
assertEquals(20, writtenFilePaths.size());
@@ -533,7 +638,7 @@ public class AddFilesIT {
p.apply(Create.of(writtenFilePaths))
.apply(
new AddFiles(
-
IcebergCatalogConfig.builder().setCatalogProperties(BIGLAKE_PROPS).build(),
+
IcebergCatalogConfig.builder().setCatalogProperties(LAKEHOUSE_PROPS).build(),
namespace + "." + destTableName,
null,
null,
@@ -581,7 +686,10 @@ public class AddFilesIT {
Managed.read(Managed.ICEBERG)
.withConfig(
ImmutableMap.of(
- "table", destTableId.toString(),
"catalog_properties", BIGLAKE_PROPS)))
+ "table",
+ destTableId.toString(),
+ "catalog_properties",
+ LAKEHOUSE_PROPS)))
.getSinglePCollection()
.apply(MapElements.into(strings()).via(AddFilesIT::canonicalRecord));
PAssert.that(destRows)
@@ -616,7 +724,10 @@ public class AddFilesIT {
Managed.read(Managed.ICEBERG)
.withConfig(
ImmutableMap.of(
- "table", destTableId.toString(),
"catalog_properties", BIGLAKE_PROPS)))
+ "table",
+ destTableId.toString(),
+ "catalog_properties",
+ LAKEHOUSE_PROPS)))
.getSinglePCollection();
PAssert.that(destRows).containsInAnyOrder(TEST_ROWS);
@@ -634,7 +745,7 @@ public class AddFilesIT {
format(
"%s.%s.%s.%s",
PROJECT,
- CATALOG_NAME,
+ LakehouseTestCatalog.CATALOG_ID,
destTableId.namespace(),
destTableId.name()))))
.getSinglePCollection()
@@ -692,7 +803,7 @@ public class AddFilesIT {
.apply(Deduplicate.values())
.apply(
new AddFiles(
-
IcebergCatalogConfig.builder().setCatalogProperties(BIGLAKE_PROPS).build(),
+
IcebergCatalogConfig.builder().setCatalogProperties(LAKEHOUSE_PROPS).build(),
namespace + "." + destTableName,
null,
PARTITION_FIELDS,
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/LakehouseTestCatalog.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/LakehouseTestCatalog.java
new file mode 100644
index 00000000000..6faac4400b3
--- /dev/null
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/LakehouseTestCatalog.java
@@ -0,0 +1,198 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
+
+import com.google.api.services.storage.model.StorageObject;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
+import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import org.apache.beam.sdk.extensions.gcp.util.GcsUtil;
+import org.apache.beam.sdk.testing.TestPipeline;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.SupportsNamespaces;
+import org.apache.iceberg.catalog.TableIdentifier;
+
+/**
+ * Test-side description of the Lakehouse Iceberg REST catalog the ITs run
against.
+ *
+ * <p>The catalog is a multiple-bucket catalog addressed by a {@code
+ * bl://projects/PROJECT/catalogs/CATALOG} warehouse and allowed to store
resources under a fixed
+ * set of Cloud Storage locations: the catalog's default location plus any
restricted locations.
+ * Configured (defaults and overrides) by the integrationTest task in
build.gradle through the
+ * system properties
+ *
+ * <ul>
+ * <li>{@code beam.iceberg.lakehouse.warehouse}: the {@code bl://} warehouse
URI
+ * <li>{@code beam.iceberg.lakehouse.locations}: comma-separated {@code
gs://} prefixes the
+ * catalog may write to; the first is the catalog's default location,
the rest are additional
+ * restricted locations
+ * </ul>
+ */
+public final class LakehouseTestCatalog {
+ private static final Pattern WAREHOUSE_PATTERN =
+ Pattern.compile("bl://projects/([^/]+)/catalogs/([^/]+)");
+
+ public static final String WAREHOUSE =
requiredProperty("beam.iceberg.lakehouse.warehouse");
+
+ public static final List<String> LOCATIONS =
+ ImmutableList.copyOf(
+ Splitter.on(',')
+ .trimResults()
+ .omitEmptyStrings()
+ .split(requiredProperty("beam.iceberg.lakehouse.locations")));
+
+ /** Catalog id, which is also the second segment of BigQuery's 4-part table
reference. */
+ public static final String CATALOG_ID = parseCatalogId(WAREHOUSE);
+
+ private static final String PROJECT =
+ TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
+
+ private LakehouseTestCatalog() {}
+
+ private static String requiredProperty(String name) {
+ String value = System.getProperty(name);
+ checkArgument(
+ value != null && !value.isEmpty(),
+ "System property %s is not set; run through the integrationTest gradle
task, which"
+ + " sets it, or pass -D%s=...",
+ name,
+ name);
+ return value;
+ }
+
+ private static String parseCatalogId(String warehouse) {
+ // Legacy single-bucket catalogs are addressed by their bucket and named
after it.
+ if (warehouse.startsWith("gs://")) {
+ return bucketOf(warehouse);
+ }
+ Matcher matcher = WAREHOUSE_PATTERN.matcher(warehouse);
+ checkArgument(
+ matcher.matches(),
+ "Expected a bl://projects/PROJECT/catalogs/CATALOG (or legacy
gs://BUCKET) warehouse, got '%s'",
+ warehouse);
+ return matcher.group(2);
+ }
+
+ /** The catalog's default location; tables land here unless created with an
explicit location. */
+ public static String defaultLocation() {
+ return LOCATIONS.get(0);
+ }
+
+ /** A restricted location outside the default one, i.e. in a second bucket.
*/
+ public static String additionalLocation() {
+ checkArgument(
+ LOCATIONS.size() >= 2,
+ "beam.iceberg.lakehouse.locations must list at least two locations,
got %s",
+ LOCATIONS);
+ return LOCATIONS.get(1);
+ }
+
+ public static String bucketOf(String gcsLocation) {
+ checkArgument(gcsLocation.startsWith("gs://"), "Not a gs:// location: %s",
gcsLocation);
+ String withoutScheme = gcsLocation.substring("gs://".length());
+ int slash = withoutScheme.indexOf('/');
+ return slash < 0 ? withoutScheme : withoutScheme.substring(0, slash);
+ }
+
+ /** Object-name prefix (no bucket, no leading slash) of a {@code
gs://bucket/path} location. */
+ public static String prefixOf(String gcsLocation) {
+ String withoutScheme = gcsLocation.substring("gs://".length());
+ int slash = withoutScheme.indexOf('/');
+ return slash < 0 ? "" : withoutScheme.substring(slash + 1);
+ }
+
+ public static Map<String, String> catalogProperties() {
+ return ImmutableMap.<String, String>builder()
+ .put("type", "rest")
+ .put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog")
+ .put("warehouse", WAREHOUSE)
+ .put("header.x-goog-user-project", PROJECT)
+ // Required by catalogs in vended-credentials mode; ignored in
end-user mode.
+ .put("header.X-Iceberg-Access-Delegation", "vended-credentials")
+ .put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO")
+ .put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager")
+ .build();
+ }
+
+ /**
+ * Drops every table in the namespaces and then the namespaces themselves,
and deletes the tables'
+ * files from Cloud Storage.
+ */
+ public static void dropNamespacesAndFiles(Catalog catalog, List<String>
namespaces)
+ throws IOException {
+ for (String name : namespaces) {
+ Namespace namespace = Namespace.of(name);
+ if (!((SupportsNamespaces) catalog).namespaceExists(namespace)) {
+ continue;
+ }
+ dropTablesAndFiles(catalog, name);
+ ((SupportsNamespaces) catalog).dropNamespace(namespace);
+ }
+ }
+
+ /**
+ * Drops every table in the namespace and deletes the tables' files from
Cloud Storage: Lakehouse
+ * keeps a dropped table's data and metadata (even with purge), and table
locations carry a random
+ * suffix, so they are captured before the drop.
+ */
+ public static void dropTablesAndFiles(Catalog catalog, String namespaceName)
throws IOException {
+ Namespace namespace = Namespace.of(namespaceName);
+ if (!((SupportsNamespaces) catalog).namespaceExists(namespace)) {
+ return;
+ }
+ List<String> tableLocations = new ArrayList<>();
+ for (TableIdentifier identifier : catalog.listTables(namespace)) {
+ tableLocations.add(catalog.loadTable(identifier).location());
+ catalog.dropTable(identifier);
+ }
+ for (String location : tableLocations) {
+ deleteObjects(location);
+ }
+ }
+
+ /** Deletes every object under a {@code gs://bucket/prefix} location. */
+ public static void deleteObjects(String gcsLocation) throws IOException {
+ GcsUtil gcsUtil =
TestPipeline.testingPipelineOptions().as(GcsOptions.class).getGcsUtil();
+ List<StorageObject> objects =
+ gcsUtil.listObjects(bucketOf(gcsLocation), prefixOf(gcsLocation),
null).getItems();
+ if (objects == null || objects.isEmpty()) {
+ return;
+ }
+ List<String> paths = new ArrayList<>();
+ for (StorageObject object : objects) {
+ paths.add("gs://" + object.getBucket() + "/" + object.getName());
+ }
+ gcsUtil.remove(paths);
+ }
+
+ /** BigQuery's 4-part {@code project.catalog.namespace.table} reference. */
+ public static String bigQueryTableSpec(String namespace, String table) {
+ return String.format("%s.%s.%s.%s", PROJECT, CATALOG_ID, namespace, table);
+ }
+}
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCatalogBaseIT.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCatalogBaseIT.java
index 859f20c8fa6..d4b8970c604 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCatalogBaseIT.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCatalogBaseIT.java
@@ -175,10 +175,9 @@ public abstract class IcebergCatalogBaseIT implements
Serializable {
/**
* Catalogs whose tables are also queryable with BigQuery return the
BigQuery table reference for
* the given Iceberg table id: either the 4-part {@code
project.catalog.namespace.table} form for
- * Lakehouse runtime catalog (BigLake metastore REST) tables, or the 3-part
{@code
- * project.dataset.table} form for the BigQuery metastore federation, where
namespaces surface as
- * datasets. Returning null (the default) disables the cross-engine read
checks in {@link
- * #testReadWithBigQueryIO()}.
+ * Lakehouse runtime catalog (Iceberg REST) tables, or the 3-part {@code
project.dataset.table}
+ * form for the BigQuery metastore federation, where namespaces surface as
datasets. Returning
+ * null (the default) disables the cross-engine read checks in {@link
#testReadWithBigQueryIO()}.
*/
public @Nullable String bigQueryTableSpec(String tableId) {
return null;
@@ -238,14 +237,14 @@ public abstract class IcebergCatalogBaseIT implements
Serializable {
try {
GcsUtil gcsUtil = OPTIONS.as(GcsOptions.class).getGcsUtil();
GcsPath path = GcsPath.fromUri(warehouse);
+ // The warehouse may be a bare bucket (no object path), where
getFileName() throws.
+ String prefix =
+ path.getObject().isEmpty()
+ ? getClass().getSimpleName()
+ : getClass().getSimpleName() + "/" + path.getFileName();
@Nullable List<StorageObject> objects =
- gcsUtil
- .listObjects(
- path.getBucket(),
- getClass().getSimpleName() + "/" +
path.getFileName().toString(),
- null)
- .getItems();
+ gcsUtil.listObjects(path.getBucket(), prefix, null).getItems();
// sometimes a catalog's cleanup will take care of all the files.
// If any files are left though, manually delete them with GCS utils
@@ -432,7 +431,7 @@ public abstract class IcebergCatalogBaseIT implements
Serializable {
}
}
- private List<Record> readRecords(Table table) throws IOException {
+ protected List<Record> readRecords(Table table) throws IOException {
org.apache.iceberg.Schema tableSchema = table.schema();
TableScan tableScan = table.newScan().project(tableSchema);
List<Record> writtenRecords = new ArrayList<>();
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCdcWriteBaseIT.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCdcWriteBaseIT.java
index f35fcd82996..d641a0e5f39 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCdcWriteBaseIT.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/IcebergCdcWriteBaseIT.java
@@ -236,13 +236,13 @@ public abstract class IcebergCdcWriteBaseIT implements
Serializable {
try {
GcsUtil gcsUtil = OPTIONS.as(GcsOptions.class).getGcsUtil();
GcsPath path = GcsPath.fromUri(warehouse);
+ // The warehouse may be a bare bucket (no object path), where
getFileName() throws.
+ String prefix =
+ path.getObject().isEmpty()
+ ? getClass().getSimpleName()
+ : getClass().getSimpleName() + "/" + path.getFileName();
@Nullable List<StorageObject> objects =
- gcsUtil
- .listObjects(
- path.getBucket(),
- getClass().getSimpleName() + "/" +
path.getFileName().toString(),
- null)
- .getItems();
+ gcsUtil.listObjects(path.getBucket(), prefix, null).getItems();
// A catalog's cleanup sometimes removes every file; delete whatever is
left.
if (objects != null) {
gcsUtil.remove(
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/LakehouseCatalogCdcWriteIT.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/LakehouseCatalogCdcWriteIT.java
index fe0d7d76451..c5e1efb0098 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/LakehouseCatalogCdcWriteIT.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/LakehouseCatalogCdcWriteIT.java
@@ -17,7 +17,9 @@
*/
package org.apache.beam.sdk.io.iceberg.catalog;
+import java.io.IOException;
import java.util.Map;
+import org.apache.beam.sdk.io.iceberg.LakehouseTestCatalog;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.rest.RESTCatalog;
@@ -28,27 +30,19 @@ import org.junit.BeforeClass;
public class LakehouseCatalogCdcWriteIT extends IcebergCdcWriteBaseIT {
private static Map<String, String> catalogProps;
- private static final String LAKEHOUSE_WAREHOUSE =
- System.getProperty("beam.iceberg.biglake.warehouse",
"gs://managed-iceberg-biglake-its");
-
@BeforeClass
public static void setup() {
- warehouse = LAKEHOUSE_WAREHOUSE;
- catalogProps =
- ImmutableMap.<String, String>builder()
- .put("type", "rest")
- .put("uri",
"https://biglake.googleapis.com/iceberg/v1/restcatalog")
- .put("warehouse", LAKEHOUSE_WAREHOUSE)
- .put("header.x-goog-user-project", OPTIONS.getProject())
- .put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO")
- .put("rest.auth.type",
"org.apache.iceberg.gcp.auth.GoogleAuthManager")
- .build();
+ warehouse = LakehouseTestCatalog.defaultLocation();
+ catalogProps = LakehouseTestCatalog.catalogProperties();
}
@After
- public void after() {
+ public void after() throws IOException {
+ // Lakehouse keeps a dropped table's files, so remove them before the base
class drops the
+ // namespace.
+ LakehouseTestCatalog.dropTablesAndFiles(catalog, namespace());
// The base class points its cleanup at this warehouse.
- warehouse = LAKEHOUSE_WAREHOUSE;
+ warehouse = LakehouseTestCatalog.defaultLocation();
}
@Override
diff --git
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java
index aa83d2b7db2..ee809b42ab8 100644
---
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java
+++
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/catalog/RESTCatalogBLMSIT.java
@@ -17,61 +17,77 @@
*/
package org.apache.beam.sdk.io.iceberg.catalog;
+import static org.apache.beam.sdk.managed.Managed.ICEBERG;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.startsWith;
+import static org.junit.Assert.assertFalse;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
import java.util.Map;
+import org.apache.beam.sdk.io.iceberg.LakehouseTestCatalog;
+import org.apache.beam.sdk.managed.Managed;
+import org.apache.beam.sdk.transforms.Create;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.iceberg.BaseTable;
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.Table;
import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.SupportsNamespaces;
import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.Record;
import org.apache.iceberg.rest.RESTCatalog;
import org.junit.After;
import org.junit.BeforeClass;
+import org.junit.Test;
-/** Tests for {@link org.apache.iceberg.rest.RESTCatalog} using BigLake
Metastore. */
+/**
+ * Tests for {@link org.apache.iceberg.rest.RESTCatalog} using a
multiple-bucket Lakehouse catalog
+ * (see {@link LakehouseTestCatalog}).
+ */
public class RESTCatalogBLMSIT extends IcebergCatalogBaseIT {
private static Map<String, String> catalogProps;
- // Using a special bucket for this test class because
- // BigLake does not support using subfolders as a warehouse (yet).
- // Overridable for local runs against a different project's catalog, e.g.
- // -Dbeam.iceberg.biglake.warehouse=gs://my-bucket (bucket-backed catalogs
are named after
- // their bucket).
- private static final String BIGLAKE_WAREHOUSE =
- System.getProperty("beam.iceberg.biglake.warehouse",
"gs://managed-iceberg-biglake-its");
-
@BeforeClass
public static void setup() {
- warehouse = BIGLAKE_WAREHOUSE;
- catalogProps =
- ImmutableMap.<String, String>builder()
- .put("type", "rest")
- .put("uri",
"https://biglake.googleapis.com/iceberg/v1/restcatalog")
- .put("warehouse", BIGLAKE_WAREHOUSE)
- .put("header.x-goog-user-project", OPTIONS.getProject())
- .put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO")
- .put("rest.auth.type",
"org.apache.iceberg.gcp.auth.GoogleAuthManager")
- .build();
+ // The catalog decides where tables go (its default location); the base
class only uses
+ // `warehouse` to sweep leftover files, so point it at that location.
+ warehouse = LakehouseTestCatalog.defaultLocation();
+ catalogProps = LakehouseTestCatalog.catalogProperties();
}
@After
public void after() {
// making sure the cleanup path is directed at the correct warehouse
- warehouse = BIGLAKE_WAREHOUSE;
+ warehouse = LakehouseTestCatalog.defaultLocation();
}
@Override
public String type() {
- return "biglake";
+ return "lakehouse";
}
@Override
public String bigQueryTableSpec(String tableId) {
- // BigQuery surfaces Lakehouse runtime catalog (BigLake metastore REST)
tables via 4-part
- // project.catalog.namespace.table identifiers; the catalog id of a
bucket-backed catalog is
- // the bucket name. Requires the caller to hold biglake.* read permissions
(e.g.
- // roles/biglake.viewer) in addition to the usual BigQuery roles.
+ // BigQuery surfaces Lakehouse runtime catalog (Iceberg REST) tables via
4-part
+ // project.catalog.namespace.table identifiers. Requires the caller to
hold biglake.* read
+ // permissions (e.g. roles/biglake.viewer) in addition to the usual
BigQuery roles.
TableIdentifier identifier = TableIdentifier.parse(tableId);
- String catalogId = BIGLAKE_WAREHOUSE.replace("gs://", "");
- return String.format(
- "%s.%s.%s.%s", OPTIONS.getProject(), catalogId,
identifier.namespace(), identifier.name());
+ return LakehouseTestCatalog.bigQueryTableSpec(
+ identifier.namespace().toString(), identifier.name());
+ }
+
+ @Override
+ public void catalogCleanup(List<Namespace> namespaces) throws IOException {
+ List<String> names = new ArrayList<>();
+ for (Namespace namespace : namespaces) {
+ names.add(namespace.toString());
+ }
+ LakehouseTestCatalog.dropNamespacesAndFiles(catalog, names);
}
@Override
@@ -88,4 +104,53 @@ public class RESTCatalogBLMSIT extends IcebergCatalogBaseIT
{
.put("catalog_properties", catalogProps)
.build();
}
+
+ /**
+ * A multiple-bucket catalog may place resources under any of its restricted
locations, not just
+ * its default one. Lakehouse pins a table under its namespace's location,
so the namespace is
+ * created in a second bucket; Beam only receives the table through the
catalog and must write
+ * wherever it was placed.
+ */
+ @Test
+ public void testWriteReadTableInAdditionalLocation() throws IOException {
+ String altNamespace = namespace() + "_alt";
+ String altTableId = altNamespace + ".test_table";
+ String namespaceLocation = LakehouseTestCatalog.additionalLocation() + "/"
+ altNamespace;
+ assertFalse(
+ "Test needs two distinct buckets",
+ LakehouseTestCatalog.bucketOf(namespaceLocation)
+
.equals(LakehouseTestCatalog.bucketOf(LakehouseTestCatalog.defaultLocation())));
+ namespacesToCleanup.add(altNamespace);
+ ((SupportsNamespaces) catalog)
+ .createNamespace(
+ Namespace.of(altNamespace), ImmutableMap.of("location",
namespaceLocation));
+ Table table = catalog.createTable(TableIdentifier.parse(altTableId),
ICEBERG_SCHEMA);
+ assertThat(table.location(), startsWith(namespaceLocation));
+
+ pipeline
+ .apply(Create.of(inputRows))
+ .setRowSchema(BEAM_SCHEMA)
+
.apply(Managed.write(ICEBERG).withConfig(managedIcebergConfig(altTableId)));
+ pipeline.run().waitUntilFinish();
+
+ table.refresh();
+ List<Record> returnedRecords = readRecords(table);
+ assertThat(
+ returnedRecords,
containsInAnyOrder(inputRows.stream().map(RECORD_FUNC::apply).toArray()));
+
+ // Both the data files Beam wrote and the metadata the catalog committed
live in the
+ // additional location, not in the catalog's default bucket.
+ List<String> dataFileLocations = new ArrayList<>();
+ for (Snapshot snapshot : table.snapshots()) {
+ for (DataFile dataFile : snapshot.addedDataFiles(table.io())) {
+ dataFileLocations.add(dataFile.location());
+ }
+ }
+ assertFalse("No data files were written", dataFileLocations.isEmpty());
+ for (String location : dataFileLocations) {
+ assertThat(location,
startsWith(LakehouseTestCatalog.additionalLocation()));
+ }
+ String metadataLocation = ((BaseTable)
table).operations().current().metadataFileLocation();
+ assertThat(metadataLocation,
startsWith(LakehouseTestCatalog.additionalLocation()));
+ }
}