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,