Class TableMetadataDriver

All Implemented Interfaces:
Serializable, HasDisplayData

@Internal public abstract class TableMetadataDriver extends PTransform<PCollection<Row>,PCollection<KV<String,@Nullable SerializableTableSpec>>>
A driver transform that extracts table identifiers from incoming Rows, deduplicates them per window, optionally bounds the cache size up to maximumCacheSize (batch pipelines only), loads their declarative metadata from the Iceberg catalog, and emits KV pairs of table identifier strings to SerializableTableSpec (or null if the table does not exist or fails to load). This is intended to be used in Beam pipelines that may utilize a large number of workers to handle Iceberg writes, where having every worker thread query for table metadata results in an excessive amount of requests and a high level of redundancy.

Can also be materialized into a broadcasted PCollectionView via asView(IcebergCatalogConfig, DynamicDestinations). By default, the cache size is uncapped. If maximumCacheSize is configured and the number of distinct tables in a window exceeds it, up to maximumCacheSize tables are sampled into the broadcasted view, while remaining destinations fall back to worker-local catalog loading. Note that maximumCacheSize is currently supported for bounded batch pipelines only.

For unbounded streaming pipelines in GlobalWindows, Deduplicate is used to deduplicate table identifiers over the configured refreshInterval (defaulting to DEFAULT_REFRESH_INTERVAL), allowing periodic refresh of table metadata when schemas evolve. Missing table signals (null specs) trigger side-input view materialization without caching the missing tables, ensuring downstream consumers are never blocked waiting for the side input.

See Also: