This is an automated email from the ASF dual-hosted git repository.

ahmedabu98 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 0291da13e0f use binding (#40455)
0291da13e0f is described below

commit 0291da13e0ff78d9dc4113b3d82ba8f1d38b744f
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Wed Oct 7 13:09:02 2026 -0700

    use binding (#40455)
---
 .../io/iceberg/cdc/SerializableChangelogTask.java  |  6 +-
 .../IcebergCdcReadSchemaTransformProviderTest.java | 71 ++++++++++++++++++++++
 2 files changed, 76 insertions(+), 1 deletion(-)

diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
index 97bcbaa5bec..7764ea2397d 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/SerializableChangelogTask.java
@@ -44,6 +44,7 @@ import org.apache.iceberg.DeletedRowsScanTask;
 import org.apache.iceberg.PartitionSpec;
 import org.apache.iceberg.Schema;
 import org.apache.iceberg.StructLike;
+import org.apache.iceberg.expressions.Binder;
 import org.apache.iceberg.expressions.Expression;
 import org.apache.iceberg.expressions.ExpressionParser;
 
@@ -162,7 +163,10 @@ public abstract class SerializableChangelogTask {
             .setSpecId(spec.specId())
             .setStart(contentScanTask.start())
             .setLength(contentScanTask.length())
-            
.setJsonExpression(ExpressionParser.toJson(contentScanTask.residual()));
+            // bound literals serialize by column type (e.g. ISO dates), which 
fromJson expects
+            .setJsonExpression(
+                ExpressionParser.toJson(
+                    Binder.bind(spec.schema().asStruct(), 
contentScanTask.residual(), false)));
 
     if (task instanceof AddedRowsScanTask) {
       AddedRowsScanTask addedRowsTask = (AddedRowsScanTask) task;
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProviderTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProviderTest.java
index 9d12389184a..3a21566c305 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProviderTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProviderTest.java
@@ -23,6 +23,10 @@ import static 
org.apache.beam.sdk.values.PCollection.IsBounded.UNBOUNDED;
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.equalTo;
 
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.time.ZoneOffset;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -71,6 +75,16 @@ public class IcebergCdcReadSchemaTransformProviderTest {
               Types.NestedField.required(4, "event_micros", 
Types.LongType.get())),
           ImmutableSet.of(1));
 
+  private static final org.apache.iceberg.Schema TEMPORAL_SCHEMA =
+      new org.apache.iceberg.Schema(
+          ImmutableList.of(
+              Types.NestedField.required(1, "id", Types.LongType.get()),
+              Types.NestedField.optional(2, "d", Types.DateType.get()),
+              Types.NestedField.optional(3, "t", Types.TimeType.get()),
+              Types.NestedField.optional(4, "ts", 
Types.TimestampType.withoutZone()),
+              Types.NestedField.optional(5, "tstz", 
Types.TimestampType.withZone())),
+          ImmutableSet.of(1));
+
   @Rule public TestDataWarehouse warehouse = new 
TestDataWarehouse(TEMPORARY_FOLDER, "default");
 
   @Rule public TestPipeline testPipeline = TestPipeline.create();
@@ -272,6 +286,63 @@ public class IcebergCdcReadSchemaTransformProviderTest {
     testPipeline.run();
   }
 
+  @Test
+  public void testManagedReadWithTemporalFilter() throws Exception {
+    String identifier = "default.table_" + 
Long.toString(UUID.randomUUID().hashCode(), 16);
+    TableIdentifier tableId = TableIdentifier.parse(identifier);
+
+    Table table = warehouse.createTable(tableId, TEMPORAL_SCHEMA);
+    LocalDate date = LocalDate.parse("2026-01-02");
+    LocalTime time = LocalTime.parse("11:00:00");
+    LocalDateTime ts = LocalDateTime.parse("2026-01-01T13:00:00");
+    // rows 2-5 each fail exactly one predicate
+    List<Record> records =
+        ImmutableList.of(
+            temporalRecord(1L, date, time, ts, ts),
+            temporalRecord(2L, date.minusDays(2), time, ts, ts),
+            temporalRecord(3L, date, time.minusHours(1), ts, ts),
+            temporalRecord(4L, date, time, ts.minusHours(2), ts),
+            temporalRecord(5L, date, time, ts, ts.minusHours(2)));
+    table
+        .newFastAppend()
+        .appendFile(warehouse.writeRecords("cdc-temporal.parquet", 
table.schema(), records))
+        .commit();
+
+    Map<String, String> properties = new HashMap<>();
+    properties.put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP);
+    properties.put("warehouse", warehouse.location);
+
+    Map<String, Object> configMap = new HashMap<>();
+    configMap.put("table", identifier);
+    configMap.put("catalog_name", "test-name");
+    configMap.put("catalog_properties", properties);
+    configMap.put("from_snapshot", table.currentSnapshot().snapshotId());
+    configMap.put("to_snapshot", table.currentSnapshot().snapshotId());
+    configMap.put(
+        "filter",
+        "d > DATE '2026-01-01' AND t > TIME '10:30:00' "
+            + "AND ts > TIMESTAMP '2026-01-01 12:00:00' "
+            + "AND tstz > TIMESTAMP '2026-01-01 12:00:00'");
+
+    Schema schema = IcebergUtils.icebergSchemaToBeamSchema(table.schema());
+    PCollection<Row> output =
+        testPipeline
+            .apply(Managed.read(Managed.ICEBERG_CDC).withConfig(configMap))
+            .getSinglePCollection();
+
+    PAssert.that(output)
+        .containsInAnyOrder(IcebergUtils.icebergRecordToBeamRow(schema, 
records.get(0)));
+
+    testPipeline.run();
+  }
+
+  private static Record temporalRecord(
+      long id, LocalDate d, LocalTime t, LocalDateTime ts, LocalDateTime tstz) 
{
+    return TestFixtures.createRecord(
+        TEMPORAL_SCHEMA,
+        ImmutableMap.of("id", id, "d", d, "t", t, "ts", ts, "tstz", 
tstz.atOffset(ZoneOffset.UTC)));
+  }
+
   private static Record record(long id, String data, String category, long 
eventMicros) {
     return TestFixtures.createRecord(
         CDC_CONFIG_SCHEMA,

Reply via email to