Class AddFiles
- All Implemented Interfaces:
Serializable,HasDisplayData
DataFile with partition metadata and column stats, batched into manifests
and committed as snapshots.
Outputs: snapshots (one row per commit), errors (one row per file that could
not be registered: file, error), and dry_run_report when a dry run is
configured.
Schema evolution. With a SchemaEvolutionConfig whose options are set, a
pre-pass reads every Parquet footer, classifies the change each distinct file schema needs on the
table (add a column, relax a required column, promote a type), commits the allowed changes in one
transaction, and only then registers the files. Manifest entries are immutable, so this ordering
is what guarantees that every registered file carries stats for every column it has. Files whose
schema needs a change that is not allowed, or that conflicts with the table or with another file,
are incompatible: by default the pipeline fails before committing anything, or routes them to
errors (see SchemaEvolutionConfig.IncompatibleSchemaHandling). The per-file
checks read Parquet footers: an ORC or Avro file, or a pinned column the footer has no null count
for, cannot be verified and goes to errors unless SchemaEvolutionConfig.UnverifiableFileHandling.ACCEPT registers it on trust. 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 errors. Schema evolution
currently requires bounded input; unbounded input with options set is rejected at construction.
SchemaEvolutionConfig evolution =
SchemaEvolutionConfig.builder()
.setOptions(EnumSet.of(ALLOW_FIELD_ADDITION, ALLOW_TYPE_PROMOTION))
.setRequiredColumns(Collections.singleton("id"))
.build();
paths.apply(new AddFiles(catalog, "db.sales", null, null, null, null, null, null, evolution));
Without options the table schema is never changed and files register as-is: columns the table does not have get no stats and are not readable, and a nested column the table does not know can make that file, and any scan that includes it, fail in an Iceberg reader.
- See Also:
-
Field Summary
Fields inherited from class org.apache.beam.sdk.transforms.PTransform
annotations, displayData, name, resourceHints -
Constructor Summary
ConstructorsConstructorDescriptionAddFiles(IcebergCatalogConfig catalogConfig, String tableIdentifier, @Nullable String locationPrefix, @Nullable List<String> partitionFields, @Nullable List<String> sortFields, @Nullable Map<String, String> tableProps, @Nullable Integer manifestFileSize, @Nullable Duration intervalTrigger) AddFiles(IcebergCatalogConfig catalogConfig, String tableIdentifier, @Nullable String locationPrefix, @Nullable List<String> partitionFields, @Nullable List<String> sortFields, @Nullable Map<String, String> tableProps, @Nullable Integer manifestFileSize, @Nullable Duration intervalTrigger, @Nullable SchemaEvolutionConfig evolution) -
Method Summary
Modifier and TypeMethodDescriptionexpand(PCollection<String> input) Override this method to specify how thisPTransformshould be expanded on the givenInputT.static org.apache.iceberg.FileFormatinferFormat(String path) Tries to infer other file formats.Methods inherited from class org.apache.beam.sdk.transforms.PTransform
addAnnotation, compose, compose, getAdditionalInputs, getAnnotations, getDefaultOutputCoder, getDefaultOutputCoder, getDefaultOutputCoder, getKindString, getName, getResourceHints, populateDisplayData, setDisplayData, setResourceHints, toString, validate, validate
-
Constructor Details
-
AddFiles
-
AddFiles
public AddFiles(IcebergCatalogConfig catalogConfig, String tableIdentifier, @Nullable String locationPrefix, @Nullable List<String> partitionFields, @Nullable List<String> sortFields, @Nullable Map<String, String> tableProps, @Nullable Integer manifestFileSize, @Nullable Duration intervalTrigger, @Nullable SchemaEvolutionConfig evolution) - Parameters:
locationPrefix- when set on a partitioned table, the partition is read from the path after this prefix instead of from the file's column statspartitionFields- partition spec applied when the table is created by this transformsortFields- sort order applied when the table is created by this transformtableProps- table properties applied when the table is created by this transformmanifestFileSize- data files per manifestintervalTrigger- streaming only: how often manifests are committedevolution- schema evolution settings; null or no options means the schema is never changed
-
-
Method Details
-
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<String>,PCollectionRowTuple>
-
inferFormat
Tries to infer other file formats. Defaults to Parquet.
-