Class IcebergIO.WriteRows

All Implemented Interfaces:
Serializable, HasDisplayData
Enclosing class:
IcebergIO

public abstract static class IcebergIO.WriteRows extends PTransform<PCollection<Row>,IcebergWriteResult>
See Also:
  • Constructor Details

    • WriteRows

      public WriteRows()
  • Method Details

    • to

      public IcebergIO.WriteRows to(org.apache.iceberg.catalog.TableIdentifier identifier)
    • to

      public IcebergIO.WriteRows to(DynamicDestinations destinations)
    • withTriggeringFrequency

      public IcebergIO.WriteRows withTriggeringFrequency(Duration triggeringFrequency)
      Sets the frequency at which data is written to files and a new Snapshot is produced.

      Roughly every triggeringFrequency duration, records are written to data files and appended to the respective table. Each append operation creates a new table snapshot.

      Generally speaking, increasing this duration will result in fewer, larger data files and fewer snapshots.

      This is only applicable when writing an unbounded PCollection (i.e. a streaming pipeline).

    • withDirectWriteByteLimit

      public IcebergIO.WriteRows withDirectWriteByteLimit(Integer directWriteByteLimit)
    • withDistributionMode

      public IcebergIO.WriteRows withDistributionMode(org.apache.iceberg.DistributionMode mode)
      Defines distribution of write data. Supported distributions:
      1. DistributionMode.NONE: don't shuffle rows (default)
      2. DistributionMode.HASH: shuffle rows by partition key before writing data
      DistributionMode.RANGE is not supported yet
    • withAutosharding

      public IcebergIO.WriteRows withAutosharding()
    • withWriteProperties

      public IcebergIO.WriteRows withWriteProperties(Map<String,String> writeProperties)
      Defines properties to be passed to the Iceberg writer itself. Note that these properties are execution-scoped, meaning that they are applied to a preexisting table and will not mutate any table-level properties.

      To set table-level properties that will be applied to dynamically created tables, use the managed Iceberg transform instead, setting the `table_properties` config property.

      See: https://iceberg.apache.org/docs/latest/configuration/#write-properties

    • withPartitionFields

      public IcebergIO.WriteRows withPartitionFields(List<String> partitionFields)
      Defines the desired Partition Spec to be applied when the Iceberg table must be dynamically created, e.g. `bucket(id_field, 32)` or `day(timestamp_field)`

      See: https://iceberg.apache.org/spec/#partitioning

    • withSortOrder

      public IcebergIO.WriteRows withSortOrder(List<String> sortFields)
      Defines the desired Sort Order to be applied when the Iceberg table must be dynamically created, e.g. `int_field desc` or `bucket(modulo_5, 4) asc nulls last`

      See: https://iceberg.apache.org/spec/#sorting

    • withSideInputTableCache

      public IcebergIO.WriteRows withSideInputTableCache()
      Enables expirable side-input caching of Iceberg table metadata across workers.

      When enabled, a driver transform periodically polls the Iceberg catalog and broadcasts lightweight table specifications as a side input. Workers construct in-memory Table representations without issuing remote catalog RPCs, drastically reducing catalog load.

    • withMaximumTableCacheSize

      public IcebergIO.WriteRows withMaximumTableCacheSize(int maximumTableCacheSize)
      Sets the maximum number of distinct table metadata specifications to broadcast in the side-input cache. Any tables exceeding this limit fall back to worker-local catalog loading.

      Note: This option is only supported for bounded (batch) pipelines. Calling this on an unbounded streaming pipeline will throw an exception at pipeline construction.

    • withTableCacheRefreshInterval

      public IcebergIO.WriteRows withTableCacheRefreshInterval(Duration refreshInterval)
      Sets the interval at which table metadata is refreshed from the Iceberg catalog.

      Applicable for unbounded streaming pipelines. Defaults to 5 minutes.

    • withTableCachePollingBuckets

      public IcebergIO.WriteRows withTableCachePollingBuckets(int pollingBuckets)
      Sets the number of parallel buckets/workers used to query the Iceberg catalog during refreshes. Defaults to 1 to serialize catalog queries and protect catalogs from connection spikes.
    • populateDisplayData

      public void populateDisplayData(DisplayData.Builder builder)
      Description copied from class: PTransform
      Register display data for the given transform or component.

      populateDisplayData(DisplayData.Builder) is invoked by Pipeline runners to collect display data via DisplayData.from(HasDisplayData). Implementations may call super.populateDisplayData(builder) in order to register display data in the current namespace, but should otherwise use subcomponent.populateDisplayData(builder) to use the namespace of the subcomponent.

      By default, does not register any display data. Implementors may override this method to provide their own display data.

      Specified by:
      populateDisplayData in interface HasDisplayData
      Overrides:
      populateDisplayData in class PTransform<PCollection<Row>,IcebergWriteResult>
      Parameters:
      builder - The builder to populate with display data.
      See Also:
    • expand

      public IcebergWriteResult expand(PCollection<Row> input)
      Description copied from class: PTransform
      Override this method to specify how this PTransform should be expanded on the given InputT.

      NOTE: This method should not be called directly. Instead apply the PTransform should be applied to the InputT using the apply method.

      Composite transforms, which are defined in terms of other transforms, should return the output of one of the composed transforms. Non-composite transforms, which do not apply any transforms internally, should return a new unbound output and register evaluators (via backend-specific registration methods).

      Specified by:
      expand in class PTransform<PCollection<Row>,IcebergWriteResult>