Class TableMetadataDriver
- All Implemented Interfaces:
Serializable,HasDisplayData
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:
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic classstatic interface -
Field Summary
FieldsFields inherited from class org.apache.beam.sdk.transforms.PTransform
annotations, displayData, name, resourceHints -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionasView()Helper that applies thisTableMetadataDriverand creates aPCollectionViewofMapof table identifier strings toSerializableTableSpec.asView(IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations) Helper that appliesTableMetadataDriverwith default configuration and creates an uncappedPCollectionViewofMapof table identifier strings toSerializableTableSpec.static TableMetadataDriver.Builderbuilder()expand(PCollection<Row> input) Override this method to specify how thisPTransformshould be expanded on the givenInputT.abstract IcebergCatalogConfigabstract DynamicDestinationsReturns the number of parallel buckets/workers used to query the Iceberg catalog, ornullfor default.voidpopulateDisplayData(DisplayData.Builder builder) Register display data for the given transform or component.abstract TableMetadataDriver.BuilderMethods inherited from class org.apache.beam.sdk.transforms.PTransform
addAnnotation, compose, compose, getAdditionalInputs, getAnnotations, getDefaultOutputCoder, getDefaultOutputCoder, getDefaultOutputCoder, getKindString, getName, getResourceHints, setDisplayData, setResourceHints, toString, validate, validate
-
Field Details
-
DEFAULT_REFRESH_INTERVAL
-
DEFAULT_POLLING_BUCKETS
public static final int DEFAULT_POLLING_BUCKETS- See Also:
-
-
Constructor Details
-
TableMetadataDriver
public TableMetadataDriver()
-
-
Method Details
-
getCatalogConfig
-
getDynamicDestinations
-
getMaximumCacheSize
-
getRefreshInterval
-
getPollingBuckets
Returns the number of parallel buckets/workers used to query the Iceberg catalog, ornullfor default. -
builder
-
toBuilder
-
asView
Helper that applies thisTableMetadataDriverand creates aPCollectionViewofMapof table identifier strings toSerializableTableSpec. -
asView
public static PTransform<PCollection<Row>,PCollectionView<Map<String, asViewSerializableTableSpec>>> (IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations) Helper that appliesTableMetadataDriverwith default configuration and creates an uncappedPCollectionViewofMapof table identifier strings toSerializableTableSpec. -
expand
Description copied from class:PTransformOverride this method to specify how thisPTransformshould be expanded on the givenInputT.NOTE: This method should not be called directly. Instead apply the
PTransformshould be applied to theInputTusing theapplymethod.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:
expandin classPTransform<PCollection<Row>,PCollection<KV<String, @Nullable SerializableTableSpec>>>
-
populateDisplayData
Description copied from class:PTransformRegister display data for the given transform or component.populateDisplayData(DisplayData.Builder)is invoked by Pipeline runners to collect display data viaDisplayData.from(HasDisplayData). Implementations may callsuper.populateDisplayData(builder)in order to register display data in the current namespace, but should otherwise usesubcomponent.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:
populateDisplayDatain interfaceHasDisplayData- Overrides:
populateDisplayDatain classPTransform<PCollection<Row>,PCollection<KV<String, @Nullable SerializableTableSpec>>> - Parameters:
builder- The builder to populate with display data.- See Also:
-