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()));
+  }
 }

Reply via email to