Class AddFiles

All Implemented Interfaces:
Serializable, HasDisplayData

public class AddFiles extends PTransform<PCollection<String>,PCollectionRowTuple>
Registers existing Parquet, ORC or Avro files in an Iceberg table without rewriting them: each path becomes a 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:
  • Constructor Details

  • Method Details

    • expand

      public PCollectionRowTuple expand(PCollection<String> input)
      Description copied from class: PTransform
      Override this method to specify how this PTransform should be expanded on the given InputT.

      NOTE: This method should not be called directly. Instead apply the PTransform should be applied to the InputT using the apply method.

      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:
      expand in class PTransform<PCollection<String>,PCollectionRowTuple>
    • inferFormat

      public static org.apache.iceberg.FileFormat inferFormat(String path)
      Tries to infer other file formats. Defaults to Parquet.