Class IcebergWriteSchemaTransformProvider.Configuration

java.lang.Object
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.Configuration
Enclosing class:
IcebergWriteSchemaTransformProvider

@DefaultSchema(AutoValueSchema.class) public abstract static class IcebergWriteSchemaTransformProvider.Configuration extends Object
  • Constructor Details

    • Configuration

      public Configuration()
  • Method Details

    • builder

    • getTable

      @SchemaFieldDescription("A fully-qualified table identifier. You may also provide a template to write to multiple dynamic destinations, for example: `dataset.my_{col1}_{col2.nested}_table`.") public abstract String getTable()
    • getCatalogName

      @SchemaFieldDescription("Name of the catalog containing the table.") public abstract @Nullable String getCatalogName()
    • getCatalogProperties

      @SchemaFieldDescription("Properties used to set up the Iceberg catalog.") public abstract @Nullable Map<String,String> getCatalogProperties()
    • getConfigProperties

      @SchemaFieldDescription("Properties passed to the Hadoop Configuration.") public abstract @Nullable Map<String,String> getConfigProperties()
    • getTriggeringFrequencySeconds

      @SchemaFieldDescription("For a streaming pipeline, sets the frequency at which snapshots are produced.") public abstract @Nullable Integer getTriggeringFrequencySeconds()
    • getDirectWriteByteLimit

      @SchemaFieldDescription("For a streaming pipeline, sets the limit for lifting bundles into the direct write path.") public abstract @Nullable Integer getDirectWriteByteLimit()
    • getMode

      @SchemaFieldDescription("Controls how rows are written. \'append\' (default) appends every row as new data. \'merge-on-read\' treats each row as a change (INSERT, UPDATE_BEFORE, UPDATE_AFTER, or DELETE) applied to the table by primary key.") public abstract @Nullable String getMode()
    • getSequenceNumberColumn

      @SchemaFieldDescription("Merge-on-read only. The required column name representing the monotonic sequence number used to order a single key\'s changes. Defaults to \'_commit_snapshot_sequence_number\'. This column will be stripped from the data row before writing to Iceberg.") public abstract @Nullable String getSequenceNumberColumn()
    • getChangeTypeColumn

      @SchemaFieldDescription("Merge-on-read only. The optional column name representing the row\'s change type (INSERT, UPDATE_BEFORE, UPDATE_AFTER, or DELETE). This column will be stripped from the data row before writing to Iceberg. If unset, the sink will use the element\'s native ValueKind") public abstract @Nullable String getChangeTypeColumn()
    • getChangeTypeMap

      @SchemaFieldDescription("Merge-on-read only. Optional map from a change_type_column value to the canonical change type name (see above).") public abstract @Nullable Map<String,String> getChangeTypeMap()
    • getUpsert

      @SchemaFieldDescription("Merge-on-read only. If true, only the after-image of each change (INSERT/UPDATE_AFTER) is applied, as an upsert; UPDATE_BEFORE records are dropped. Default: false.") public abstract @Nullable Boolean getUpsert()
    • getKeep

      @SchemaFieldDescription("A list of field names to keep in the input record. All other fields are dropped before writing. Is mutually exclusive with \'drop\' and \'only\'. In merge-on-read mode the control columns are dropped unless listed here.") public abstract @Nullable List<String> getKeep()
    • getDrop

      @SchemaFieldDescription("A list of field names to drop from the input record before writing. Is mutually exclusive with \'keep\' and \'only\'. In merge-on-read mode the control columns are always dropped.") public abstract @Nullable List<String> getDrop()
    • getOnly

      @SchemaFieldDescription("The name of a single record field that should be written. Is mutually exclusive with \'keep\' and \'drop\'.") public abstract @Nullable String getOnly()
    • getPartitionFields

      @SchemaFieldDescription("Fields used to create a partition spec that is applied when tables are created. For a field \'foo\', the available partition transforms are:\n\n- `foo`\n- `truncate(foo, N)`\n- `bucket(foo, N)`\n- `hour(foo)`\n- `day(foo)`\n- `month(foo)`\n- `year(foo)`\n- `void(foo)`\n\nFor more information on partition transforms, please visit https://iceberg.apache.org/spec/#partition-transforms.") public abstract @Nullable List<String> getPartitionFields()
    • getTableProperties

      @SchemaFieldDescription("Iceberg table properties to be set on the table when it is created.\nFor more information on table properties, please visit https://iceberg.apache.org/docs/latest/configuration/#table-properties.") public abstract @Nullable Map<String,String> getTableProperties()
    • getSortFields

      @SchemaFieldDescription("Fields used to set the table\'s sort order, applied when the table is created. Each entry has the form `<term> [asc|desc] [nulls first|nulls last]`, where `<term>` is a field name or one of the partition transforms (e.g. `bucket(col, 4)`, `day(ts)`). Direction defaults to ascending; null order defaults to nulls-first for ascending and nulls-last for descending. Note: this sets the table\'s declared sort order as metadata; it does not cause Beam to physically sort records before writing.\nFor more information on sort orders, please visit https://iceberg.apache.org/spec/#sort-orders.") public abstract @Nullable List<String> getSortFields()
    • getDistributionMode

      @SchemaFieldDescription("Defines distribution of write data. Supported distributions:\n- none: don\'t shuffle rows (default)\n- hash: shuffle rows by partition key before writing data") public abstract @Nullable String getDistributionMode()
    • getAutosharding

      @SchemaFieldDescription("Enables dynamic sharding to automatically adjust the number of parallel writers based on data volume. It handles data skew by further sub-dividing partitions into multiple shards to prevent bottlenecks during high-throughput writes. Only available with \'hash\' distribution mode.") public abstract @Nullable Boolean getAutosharding()
    • getWriteProperties

      @SchemaFieldDescription("Properties applied to the underlying file writer (e.g. Parquet write properties like \'write.parquet.bloom-filter-enabled.column.<col>\').") public abstract @Nullable Map<String,String> getWriteProperties()
    • getUseSideInputTableCache

      @SchemaFieldDescription("Enables expirable side-input caching of Iceberg table metadata across workers to reduce catalog load.") public abstract @Nullable Boolean getUseSideInputTableCache()
    • getTableCacheRefreshIntervalSeconds

      @SchemaFieldDescription("For a streaming pipeline, sets the interval in seconds at which table metadata is refreshed from the catalog.") public abstract @Nullable Integer getTableCacheRefreshIntervalSeconds()
    • getMaximumTableCacheSize

      @SchemaFieldDescription("For a batch pipeline, sets the maximum number of table metadata specs to cache in memory. Tables exceeding this limit fall back to worker-local catalog loading.") public abstract @Nullable Integer getMaximumTableCacheSize()
    • getTableCachePollingBuckets

      @SchemaFieldDescription("Sets the number of parallel buckets/workers used to query the Iceberg catalog during refreshes. Defaults to 1.") public abstract @Nullable Integer getTableCachePollingBuckets()
    • getEqualityColumns

      @SchemaFieldDescription("Columns defining row identity (equality-delete fields). Defaults to the destination table\'s identifier (primary-key) fields. Required if the table doesn\'t exist yet. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable List<String> getEqualityColumns()
    • getNumShards

      @SchemaFieldDescription("The number of deterministic primary-key-hash shards per destination, i.e. the max write parallelism per destination. Too low may bottleneck writes, and too high may produce more files. Defaults to 16. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Integer getNumShards()
    • getShardsPerPartition

      @SchemaFieldDescription("Maximum number of shards a single partition\'s rows may occupy. Lower values write fewer files per commit, but also reduces per-partition write parallelism. A value of 1 pins each partition to one writer. Ignored for unpartitioned tables. Must be between 1 and `num_shards`; defaults to `num_shards`. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Integer getShardsPerPartition()
    • getAllowedLatenessSeconds

      @SchemaFieldDescription("How long a late record may lag behind the watermark before it is dropped entirely, rather than routed to the dead_letter output. Defaults to 21600 (6 hours). Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Integer getAllowedLatenessSeconds()
    • getSinkId

      @SchemaFieldDescription("A stable identifier for this sink, used to namespace the idempotency tokens written to each commit\'s Iceberg snapshot summary. Defaults to a unique per-write UUID. Set it explicitly (and keep it stable across relaunches) for exactly-once commits across relaunches of a particular streaming write. A batch load with a stable sink_id commits only once (later batch loads with the same sink_id are skipped). Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable String getSinkId()
    • getTokenHeartbeatSeconds

      @SchemaFieldDescription("Streaming only. If set, the sink will emit a periodic empty token-refresh commit while idle, so its thread of `sink_id` stamped snapshot stays recent and is less likely to be lost to `expire_snapshots`. Disabled by default. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Integer getTokenHeartbeatSeconds()
    • getSnapshotProperties

      @SchemaFieldDescription("Extra key/value properties to add to every commit\'s Iceberg snapshot summary. Keys prefixed with \'beam.cdc.\' are reserved and rejected. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Map<String,String> getSnapshotProperties()
    • getErrorHandling

      @SchemaFieldDescription("Whether and where to output per-record invalid rows (null or missing sequence value, unknown change type, null equality value, unresolvable destination). Fails the pipeline if unset (default). Distinct from the `dead_letter` output, which is for late-but-valid rows. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable ErrorHandling getErrorHandling()
    • getSorterMemoryMb

      @SchemaFieldDescription("The in-memory buffer size (MB) for the pre-write sort; groups larger than this spill to disk. Must be >= 1. Defaults to 100. Currently only supported in \'merge-on-read\' mode.") public abstract @Nullable Integer getSorterMemoryMb()
    • getIcebergCatalog

      public IcebergCatalogConfig getIcebergCatalog()