Package org.apache.iceberg
Class BeamBaseIncrementalChangelogScan
java.lang.Object
org.apache.iceberg.BeamBaseIncrementalChangelogScan
- All Implemented Interfaces:
org.apache.iceberg.IncrementalChangelogScan,org.apache.iceberg.IncrementalScan<org.apache.iceberg.IncrementalChangelogScan,,org.apache.iceberg.ChangelogScanTask, org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>> org.apache.iceberg.Scan<org.apache.iceberg.IncrementalChangelogScan,org.apache.iceberg.ChangelogScanTask, org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>>
public class BeamBaseIncrementalChangelogScan
extends Object
implements org.apache.iceberg.IncrementalChangelogScan
Copied over from Iceberg PR #14264.
-
Field Summary
FieldsModifier and TypeFieldDescriptionprotected static final boolean -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionorg.apache.iceberg.IncrementalChangelogScancaseSensitive(boolean arg0) protected org.apache.iceberg.TableScanContextcontext()protected org.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ChangelogScanTask> doPlanFiles(Long fromSnapshotIdExclusive, long toSnapshotIdInclusive) Supplier<org.apache.iceberg.io.FileIO> fileIO()org.apache.iceberg.expressions.Expressionfilter()org.apache.iceberg.IncrementalChangelogScanfilter(org.apache.iceberg.expressions.Expression arg0) org.apache.iceberg.IncrementalChangelogScanfromSnapshotExclusive(long arg0) org.apache.iceberg.IncrementalChangelogScanfromSnapshotExclusive(String arg0) org.apache.iceberg.IncrementalChangelogScanfromSnapshotInclusive(long arg0) org.apache.iceberg.IncrementalChangelogScanfromSnapshotInclusive(String arg0) org.apache.iceberg.IncrementalChangelogScanorg.apache.iceberg.IncrementalChangelogScanorg.apache.iceberg.IncrementalChangelogScanincludeColumnStats(Collection<String> arg0) protected org.apache.iceberg.io.FileIOio()Deprecated.booleanorg.apache.iceberg.IncrementalChangelogScanmetricsReporter(org.apache.iceberg.metrics.MetricsReporter arg0) org.apache.iceberg.IncrementalChangelogScanminRowsRequested(long arg0) protected org.apache.iceberg.IncrementalChangelogScannewRefinedScan(org.apache.iceberg.Table newTable, org.apache.iceberg.Schema newSchema, org.apache.iceberg.TableScanContext newContext) org.apache.iceberg.IncrementalChangelogScanoptions()protected ExecutorServiceorg.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ChangelogScanTask> org.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>> org.apache.iceberg.IncrementalChangelogScanplanWith(ExecutorService arg0) org.apache.iceberg.IncrementalChangelogScanproject(org.apache.iceberg.Schema arg0) protected org.apache.iceberg.expressions.Expressionorg.apache.iceberg.Schemaschema()schemas()org.apache.iceberg.IncrementalChangelogScanselect(Collection<String> arg0) protected booleanprotected booleanprotected booleanintlongorg.apache.iceberg.Tabletable()protected org.apache.iceberg.Schemalongorg.apache.iceberg.IncrementalChangelogScantoSnapshot(long arg0) org.apache.iceberg.IncrementalChangelogScantoSnapshot(String arg0) org.apache.iceberg.IncrementalChangelogScanMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.iceberg.IncrementalScan
fromSnapshotExclusive, fromSnapshotExclusive, fromSnapshotInclusive, fromSnapshotInclusive, toSnapshot, toSnapshot, useBranchMethods inherited from interface org.apache.iceberg.Scan
caseSensitive, fileIO, filter, filter, ignoreResiduals, includeColumnStats, includeColumnStats, isCaseSensitive, metricsReporter, minRowsRequested, option, planFiles, planWith, project, schema, select, select, splitLookback, splitOpenFileCost, targetSplitSize
-
Field Details
-
SCAN_COLUMNS
-
SCAN_WITH_STATS_COLUMNS
-
DELETE_SCAN_COLUMNS
-
DELETE_SCAN_WITH_STATS_COLUMNS
-
PLAN_SCANS_WITH_WORKER_POOL
protected static final boolean PLAN_SCANS_WITH_WORKER_POOL
-
-
Constructor Details
-
BeamBaseIncrementalChangelogScan
public BeamBaseIncrementalChangelogScan(org.apache.iceberg.Table table)
-
-
Method Details
-
newRefinedScan
protected org.apache.iceberg.IncrementalChangelogScan newRefinedScan(org.apache.iceberg.Table newTable, org.apache.iceberg.Schema newSchema, org.apache.iceberg.TableScanContext newContext) -
doPlanFiles
protected org.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ChangelogScanTask> doPlanFiles(Long fromSnapshotIdExclusive, long toSnapshotIdInclusive) -
planTasks
public org.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>> planTasks()- Specified by:
planTasksin interfaceorg.apache.iceberg.Scan<org.apache.iceberg.IncrementalChangelogScan,org.apache.iceberg.ChangelogScanTask, org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>>
-
fromSnapshotInclusive
- Specified by:
fromSnapshotInclusivein interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
fromSnapshotInclusive
public org.apache.iceberg.IncrementalChangelogScan fromSnapshotInclusive(long arg0) - Specified by:
fromSnapshotInclusivein interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
fromSnapshotExclusive
- Specified by:
fromSnapshotExclusivein interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
fromSnapshotExclusive
public org.apache.iceberg.IncrementalChangelogScan fromSnapshotExclusive(long arg0) - Specified by:
fromSnapshotExclusivein interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
toSnapshot
public org.apache.iceberg.IncrementalChangelogScan toSnapshot(long arg0) - Specified by:
toSnapshotin interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
toSnapshot
- Specified by:
toSnapshotin interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
useBranch
- Specified by:
useBranchin interfaceorg.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
planFiles
public org.apache.iceberg.io.CloseableIterable<org.apache.iceberg.ChangelogScanTask> planFiles()- Specified by:
planFilesin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
table
public org.apache.iceberg.Table table() -
io
Deprecated. -
fileIO
- Specified by:
fileIOin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
tableSchema
protected org.apache.iceberg.Schema tableSchema() -
schemas
-
context
protected org.apache.iceberg.TableScanContext context() -
options
-
scanColumns
-
shouldReturnColumnStats
protected boolean shouldReturnColumnStats() -
columnsToKeepStats
-
shouldIgnoreResiduals
protected boolean shouldIgnoreResiduals() -
residualFilter
protected org.apache.iceberg.expressions.Expression residualFilter() -
shouldPlanWithExecutor
protected boolean shouldPlanWithExecutor() -
planExecutor
-
option
- Specified by:
optionin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
project
public org.apache.iceberg.IncrementalChangelogScan project(org.apache.iceberg.Schema arg0) - Specified by:
projectin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
caseSensitive
public org.apache.iceberg.IncrementalChangelogScan caseSensitive(boolean arg0) - Specified by:
caseSensitivein interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
isCaseSensitive
public boolean isCaseSensitive()- Specified by:
isCaseSensitivein interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
includeColumnStats
public org.apache.iceberg.IncrementalChangelogScan includeColumnStats()- Specified by:
includeColumnStatsin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
includeColumnStats
- Specified by:
includeColumnStatsin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
select
- Specified by:
selectin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
filter
public org.apache.iceberg.IncrementalChangelogScan filter(org.apache.iceberg.expressions.Expression arg0) - Specified by:
filterin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
filter
public org.apache.iceberg.expressions.Expression filter()- Specified by:
filterin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
ignoreResiduals
public org.apache.iceberg.IncrementalChangelogScan ignoreResiduals()- Specified by:
ignoreResidualsin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
planWith
- Specified by:
planWithin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
schema
public org.apache.iceberg.Schema schema()- Specified by:
schemain interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
targetSplitSize
public long targetSplitSize()- Specified by:
targetSplitSizein interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
splitLookback
public int splitLookback()- Specified by:
splitLookbackin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
splitOpenFileCost
public long splitOpenFileCost()- Specified by:
splitOpenFileCostin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
metricsReporter
public org.apache.iceberg.IncrementalChangelogScan metricsReporter(org.apache.iceberg.metrics.MetricsReporter arg0) - Specified by:
metricsReporterin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-
minRowsRequested
public org.apache.iceberg.IncrementalChangelogScan minRowsRequested(long arg0) - Specified by:
minRowsRequestedin interfaceorg.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask, G extends org.apache.iceberg.ScanTaskGroup<T>>
-