Class SchemaEvolutionConfig

java.lang.Object
org.apache.beam.sdk.io.iceberg.SchemaEvolutionConfig
All Implemented Interfaces:
Serializable

public abstract class SchemaEvolutionConfig extends Object implements Serializable
Schema evolution settings for AddFiles. With no options the table schema is never changed and files register as on a plain AddFiles; every other setting requires at least one option.

 SchemaEvolutionConfig.builder()
     .setOptions(EnumSet.of(ALLOW_FIELD_ADDITION, ALLOW_FIELD_RELAXATION, ALLOW_TYPE_PROMOTION))
     .setRequiredColumns(Set.of("id", "address.city"))   // never relaxed
     .setIncompatibleSchemaHandling(IncompatibleSchemaHandling.ROUTE_TO_ERRORS)
     .build();
 

Pins. Required columns are pinned: never made optional whatever the options say, and created required when this transform creates the table. A Parquet file that lacks a pinned column or has nulls in it is routed to the error output. Pins name canonical (table) paths, dotted for nested fields, with the container segment spelled out under lists and maps ( addresses.element.city, attributes.value.total). A top-level column whose own name contains a dot cannot be pinned.

Unverifiable files. The per-file checks read Parquet footers. An ORC or Avro file cannot be checked at all, and a Parquet file whose footer carries no null-count statistics for a pinned column (a writer with statistics disabled, or a pin under a list or map, whose physical chunk path the check does not map) cannot prove the pin. SchemaEvolutionConfig.UnverifiableFileHandling decides whether such a file is routed to the error output (the default) or registered on trust.

Incompatible schemas. A schema that needs a change the options do not allow, or that conflicts with the table or with another file's schema. SchemaEvolutionConfig.IncompatibleSchemaHandling decides whether that fails the pipeline before any schema commit (the batch default) or skips the schema so its files reach the error output (the streaming default). Files whose footer cannot be read or converted always go to the error output and never fail the pipeline. When the table does not exist, the pre-pass creates it from the union of the file schemas; if no readable Parquet schema can seed it, nothing is created and every file goes to the error output.

Dry run. Reports what a real run would do on the dry_run_report output, one row per window; nothing is committed or registered. allowed is true when every file schema can be merged and the configuration raises no problem; otherwise reason says what a real run would do about it. schemas holds one entry per distinct file schema with the changes a real run would make and, when the schema cannot be merged, the option or conflict to fix; created_table shows the table a real run would create. The report is a PCollection like any other, so attach a sink to keep it; it is also logged at INFO and the file counts are published as counters (numDryRunFilesAllowed, numDryRunFilesIncompatible, numDryRunFilesUnreadable, numDryRunFilesUnchecked, numDryRunConfigProblems). Adjust the settings, rerun until the report is allowed, then run for real with an error output attached. Against a missing table the dry run computes the union through the catalog's create-transaction API, which a REST catalog serves as a stage-create request: the credentials need table-create permission even though no table is created.

See Also: