diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ChangelogScanner.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ChangelogScanner.java index 0000884c90a5..dd24725cffbc 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ChangelogScanner.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ChangelogScanner.java @@ -51,7 +51,6 @@ import org.apache.beam.sdk.values.TupleTag; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; import org.apache.iceberg.AddedRowsScanTask; -import org.apache.iceberg.BaseIncrementalChangelogScan; import org.apache.iceberg.ChangelogScanTask; import org.apache.iceberg.DataFile; import org.apache.iceberg.DataOperations; @@ -227,10 +226,9 @@ public void process(@Element Long snapshotId, MultiOutputReceiver out) throws IO @Nullable Long fromSnapshotId = snapshot.parentId(); @Nullable Expression filter = scanConfig.getFilter(); - // TODO(ahmedabu98): replace this with table.newIncrementalChangelogScan() when - // https://github.com/apache/iceberg/pull/14264/ gets merged and released. IncrementalChangelogScan scan = - new BaseIncrementalChangelogScan(table) + table + .newIncrementalChangelogScan() .toSnapshot(snapshotId) .project(scanConfig.getProjectedSchema()); if (fromSnapshotId != null) {