ahmedabu98 commented on code in PR #40140:
URL: https://github.com/apache/beam/pull/40140#discussion_r4039485741


##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java:
##########
@@ -549,6 +630,26 @@ public IcebergWriteResult expand(PCollection<Row> input) {
             "Must only provide direct write limit for unbounded pipelines.");
       }
 
+      PCollectionView<Map<String, SerializableTableSpec>> metadataView = null;
+      if (getUsingSideInputTableCache()) {

Review Comment:
   Can we also throw if the sub-options are set without side-input enabled?



##########
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java:
##########
@@ -521,4 +527,24 @@ public void testDisplayData() {
     assertEquals("600000", items.get("tableRefreshInterval"));
     assertEquals("3", items.get("pollingBuckets"));
   }
+
+  private static class EvolveSchemaMidExecutionDoFn extends DoFn<Row, Row> {
+    private final IcebergCatalogConfig catalogConfig;
+    private final String tableIdString;
+
+    EvolveSchemaMidExecutionDoFn(IcebergCatalogConfig catalogConfig, String 
tableIdString) {
+      this.catalogConfig = catalogConfig;
+      this.tableIdString = tableIdString;
+    }
+
+    @ProcessElement
+    public void processElement(@Element Row row, OutputReceiver<Row> out) {
+      Table table =
+          
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(tableIdString));
+      if (table.schema().findField("city") == null) {
+        table.updateSchema().addColumn("city", 
Types.StringType.get()).commit();
+      }

Review Comment:
   I'm realizing we can't really use schema evolution to prove this because 
PCollection schemas are fixed. We won't be able to test with two different rows
   
   We can use partition spec evolution (e.g. 
`table.updateSpec().addField("city").commit();`). And validate after the 
pipeline runs by loading the table and verifying that:
   1. two DataFiles were written (one for each row)
   2. the first DataFile has the old spec (unpartitioned)
   3. the second DataFile has the new spec (e.g. "city")
   
   Can use this to check added files:
   `SnapshotChanges.builderFor(table).build().addedDataFiles();`



##########
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java:
##########
@@ -443,7 +443,13 @@ public void 
testStreamingSchemaEvolutionWithoutPipelineRestart() throws Exceptio
             .advanceProcessingTime(Duration.standardSeconds(3))
             .advanceWatermarkToInfinity();
 
-    PCollection<Row> input = testPipeline.apply("StreamingEvolvedInput", 
testStream);
+    PCollection<Row> input =
+        testPipeline
+            .apply("StreamingEvolvedInput", testStream)
+            .apply(
+                "EvolveSchemaMidExecution",
+                ParDo.of(new EvolveSchemaMidExecutionDoFn(catalogConfig, 
tableId.toString())))
+            .setRowSchema(v2BeamSchema);

Review Comment:
   This updates the table before the record reaches the sink or the metadata 
view step, so TableMetadataDriver's first poll is still after the update 
happens.
   
   We'd need an initial record to pass through, then a second record that 
triggers the update.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to