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..17ab4c5d30cc 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,7 @@ 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.BeamBaseIncrementalChangelogScan; import org.apache.iceberg.ChangelogScanTask; import org.apache.iceberg.DataFile; import org.apache.iceberg.DataOperations; @@ -230,7 +230,7 @@ public void process(@Element Long snapshotId, MultiOutputReceiver out) throws IO // TODO(ahmedabu98): replace this with table.newIncrementalChangelogScan() when // https://github.com/apache/iceberg/pull/14264/ gets merged and released. IncrementalChangelogScan scan = - new BaseIncrementalChangelogScan(table) + new BeamBaseIncrementalChangelogScan(table) .toSnapshot(snapshotId) .project(scanConfig.getProjectedSchema()); if (fromSnapshotId != null) { diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BaseIncrementalChangelogScan.java b/sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BeamBaseIncrementalChangelogScan.java similarity index 99% rename from sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BaseIncrementalChangelogScan.java rename to sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BeamBaseIncrementalChangelogScan.java index 6b12e16690ba..a7fe19218c1e 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BaseIncrementalChangelogScan.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BeamBaseIncrementalChangelogScan.java @@ -58,7 +58,7 @@ * Copied over from Iceberg PR #14264. */ @SuppressWarnings("nullness") -public class BaseIncrementalChangelogScan +public class BeamBaseIncrementalChangelogScan extends BaseIncrementalScan< IncrementalChangelogScan, ChangelogScanTask, ScanTaskGroup> implements IncrementalChangelogScan { @@ -80,20 +80,20 @@ private static DeleteFileIndex createEmptyInstance() { } } - private static final Logger LOG = LoggerFactory.getLogger(BaseIncrementalChangelogScan.class); + private static final Logger LOG = LoggerFactory.getLogger(BeamBaseIncrementalChangelogScan.class); - public BaseIncrementalChangelogScan(Table table) { + public BeamBaseIncrementalChangelogScan(Table table) { this(table, table.schema(), TableScanContext.empty()); } - private BaseIncrementalChangelogScan(Table table, Schema schema, TableScanContext context) { + private BeamBaseIncrementalChangelogScan(Table table, Schema schema, TableScanContext context) { super(table, schema, context); } @Override protected IncrementalChangelogScan newRefinedScan( Table newTable, Schema newSchema, TableScanContext newContext) { - return new BaseIncrementalChangelogScan(newTable, newSchema, newContext); + return new BeamBaseIncrementalChangelogScan(newTable, newSchema, newContext); } // Private fields to track build call count and cache (accessed via package-private methods for