Package org.apache.beam.sdk.io.iceberg
package org.apache.beam.sdk.io.iceberg
Iceberg connectors.
-
ClassDescriptionRegisters existing Parquet, ORC or Avro files in an Iceberg table without rewriting them: each path becomes a
DataFilewith partition metadata and column stats, batched into manifests and committed as snapshots.A wrapper that adapts a BeamRowto Iceberg'sStructLikeinterface.BundleLifter<T>A PTransform that buffers elements and outputs them to one of two TupleTags based on the total size of the bundle in finish_bundle.Utilities that convert between a SQL filter expression and an IcebergExpression.SchemaTransform implementation forIcebergIO.readRows(org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig).A connector that reads and writes to Apache Iceberg tables.SchemaTransform implementation forIcebergIO.readRows(org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig).Utilities for converting between Beam and Iceberg types, made public for user's convenience.The output of anIcebergIOwrite: the snapshots each destination table committed, plus the two diversion outputs the CDC sink can produce.SchemaTransform implementation forIcebergIO.writeRows(org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig)and, inmerge-on-readmode,IcebergIO.writeCdcRows(org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig).Helper class for source operations.Schema evolution settings forAddFiles.What to do with a file the per-file checks cannot verify: a non-Parquet file, or a Parquet file with no null-count statistics for a pinned column.Kinds of table schema changeAddFilesmay make so that a file's columns are covered by the table schema.Serializable version of an IcebergDataFile.A serializable, lightweight representation of an IcebergTable's declarative metadata.A lightweight adapter that implementsTablebacked by aSerializableTableSpec.This is an AutoValue representation of an IcebergSnapshot.Process-wide cache for IcebergTables.A driver transform that extracts table identifiers from incomingRows, deduplicates them per window, optionally bounds the cache size up tomaximumCacheSize(batch pipelines only), loads their declarative metadata from the Iceberg catalog, and emitsKVpairs of table identifier strings toSerializableTableSpec(ornullif the table does not exist or fails to load).