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]