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 Details

    • SCAN_COLUMNS

      protected static final List<String> SCAN_COLUMNS
    • SCAN_WITH_STATS_COLUMNS

      protected static final List<String> SCAN_WITH_STATS_COLUMNS
    • DELETE_SCAN_COLUMNS

      protected static final List<String> DELETE_SCAN_COLUMNS
    • DELETE_SCAN_WITH_STATS_COLUMNS

      protected static final List<String> 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:
      planTasks in interface org.apache.iceberg.Scan<org.apache.iceberg.IncrementalChangelogScan,org.apache.iceberg.ChangelogScanTask,org.apache.iceberg.ScanTaskGroup<org.apache.iceberg.ChangelogScanTask>>
    • fromSnapshotInclusive

      public org.apache.iceberg.IncrementalChangelogScan fromSnapshotInclusive(String arg0)
      Specified by:
      fromSnapshotInclusive in interface org.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:
      fromSnapshotInclusive in interface org.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • fromSnapshotExclusive

      public org.apache.iceberg.IncrementalChangelogScan fromSnapshotExclusive(String arg0)
      Specified by:
      fromSnapshotExclusive in interface org.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:
      fromSnapshotExclusive in interface org.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:
      toSnapshot in interface org.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • toSnapshot

      public org.apache.iceberg.IncrementalChangelogScan toSnapshot(String arg0)
      Specified by:
      toSnapshot in interface org.apache.iceberg.IncrementalScan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • useBranch

      public org.apache.iceberg.IncrementalChangelogScan useBranch(String arg0)
      Specified by:
      useBranch in interface org.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:
      planFiles in interface org.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 protected org.apache.iceberg.io.FileIO io()
      Deprecated.
    • fileIO

      public Supplier<org.apache.iceberg.io.FileIO> fileIO()
      Specified by:
      fileIO in interface org.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

      protected Map<Integer,org.apache.iceberg.Schema> schemas()
    • context

      protected org.apache.iceberg.TableScanContext context()
    • options

      protected Map<String,String> options()
    • scanColumns

      protected List<String> scanColumns()
    • shouldReturnColumnStats

      protected boolean shouldReturnColumnStats()
    • columnsToKeepStats

      protected Set<Integer> columnsToKeepStats()
    • shouldIgnoreResiduals

      protected boolean shouldIgnoreResiduals()
    • residualFilter

      protected org.apache.iceberg.expressions.Expression residualFilter()
    • shouldPlanWithExecutor

      protected boolean shouldPlanWithExecutor()
    • planExecutor

      protected ExecutorService planExecutor()
    • option

      public org.apache.iceberg.IncrementalChangelogScan option(String arg0, String arg1)
      Specified by:
      option in interface org.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:
      project in interface org.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:
      caseSensitive in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • isCaseSensitive

      public boolean isCaseSensitive()
      Specified by:
      isCaseSensitive in interface org.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:
      includeColumnStats in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • includeColumnStats

      public org.apache.iceberg.IncrementalChangelogScan includeColumnStats(Collection<String> arg0)
      Specified by:
      includeColumnStats in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • select

      public org.apache.iceberg.IncrementalChangelogScan select(Collection<String> arg0)
      Specified by:
      select in interface org.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:
      filter in interface org.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:
      filter in interface org.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:
      ignoreResiduals in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • planWith

      public org.apache.iceberg.IncrementalChangelogScan planWith(ExecutorService arg0)
      Specified by:
      planWith in interface org.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:
      schema in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • targetSplitSize

      public long targetSplitSize()
      Specified by:
      targetSplitSize in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • splitLookback

      public int splitLookback()
      Specified by:
      splitLookback in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>
    • splitOpenFileCost

      public long splitOpenFileCost()
      Specified by:
      splitOpenFileCost in interface org.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:
      metricsReporter in interface org.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:
      minRowsRequested in interface org.apache.iceberg.Scan<ThisT,T extends org.apache.iceberg.ScanTask,G extends org.apache.iceberg.ScanTaskGroup<T>>