Skip to content

Commit 5a3bbba

Browse files
authored
rename scanner (#39304)
1 parent f20fca2 commit 5a3bbba

2 files changed

Lines changed: 7 additions & 7 deletions

File tree

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/ChangelogScanner.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,7 @@
5151
import org.apache.beam.sdk.values.TupleTag;
5252
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
5353
import org.apache.iceberg.AddedRowsScanTask;
54-
import org.apache.iceberg.BaseIncrementalChangelogScan;
54+
import org.apache.iceberg.BeamBaseIncrementalChangelogScan;
5555
import org.apache.iceberg.ChangelogScanTask;
5656
import org.apache.iceberg.DataFile;
5757
import org.apache.iceberg.DataOperations;
@@ -230,7 +230,7 @@ public void process(@Element Long snapshotId, MultiOutputReceiver out) throws IO
230230
// TODO(ahmedabu98): replace this with table.newIncrementalChangelogScan() when
231231
// https://github.com/apache/iceberg/pull/14264/ gets merged and released.
232232
IncrementalChangelogScan scan =
233-
new BaseIncrementalChangelogScan(table)
233+
new BeamBaseIncrementalChangelogScan(table)
234234
.toSnapshot(snapshotId)
235235
.project(scanConfig.getProjectedSchema());
236236
if (fromSnapshotId != null) {

sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BaseIncrementalChangelogScan.java renamed to sdks/java/io/iceberg/src/main/java/org/apache/iceberg/BeamBaseIncrementalChangelogScan.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@
5858
* Copied over from <a href="https://github.com/apache/iceberg/pull/14264/">Iceberg PR #14264</a>.
5959
*/
6060
@SuppressWarnings("nullness")
61-
public class BaseIncrementalChangelogScan
61+
public class BeamBaseIncrementalChangelogScan
6262
extends BaseIncrementalScan<
6363
IncrementalChangelogScan, ChangelogScanTask, ScanTaskGroup<ChangelogScanTask>>
6464
implements IncrementalChangelogScan {
@@ -80,20 +80,20 @@ private static DeleteFileIndex createEmptyInstance() {
8080
}
8181
}
8282

83-
private static final Logger LOG = LoggerFactory.getLogger(BaseIncrementalChangelogScan.class);
83+
private static final Logger LOG = LoggerFactory.getLogger(BeamBaseIncrementalChangelogScan.class);
8484

85-
public BaseIncrementalChangelogScan(Table table) {
85+
public BeamBaseIncrementalChangelogScan(Table table) {
8686
this(table, table.schema(), TableScanContext.empty());
8787
}
8888

89-
private BaseIncrementalChangelogScan(Table table, Schema schema, TableScanContext context) {
89+
private BeamBaseIncrementalChangelogScan(Table table, Schema schema, TableScanContext context) {
9090
super(table, schema, context);
9191
}
9292

9393
@Override
9494
protected IncrementalChangelogScan newRefinedScan(
9595
Table newTable, Schema newSchema, TableScanContext newContext) {
96-
return new BaseIncrementalChangelogScan(newTable, newSchema, newContext);
96+
return new BeamBaseIncrementalChangelogScan(newTable, newSchema, newContext);
9797
}
9898

9999
// Private fields to track build call count and cache (accessed via package-private methods for

0 commit comments

Comments
 (0)