Package org.apache.beam.sdk.io.iceberg
Class AddFilesSchemaTransformProvider.Configuration
java.lang.Object
org.apache.beam.sdk.io.iceberg.AddFilesSchemaTransformProvider.Configuration
- Enclosing class:
AddFilesSchemaTransformProvider
@DefaultSchema(AutoValueSchema.class)
public abstract static class AddFilesSchemaTransformProvider.Configuration
extends Object
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic class -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionbuilder()abstract @Nullable ErrorHandlingValidates and converts the schema evolution settings; null when none are set.abstract StringgetTable()
-
Constructor Details
-
Configuration
public Configuration()
-
-
Method Details
-
builder
-
getTable
-
getCatalogProperties
@SchemaFieldDescription("Properties used to set up the Iceberg catalog.") public abstract @Nullable Map<String,String> getCatalogProperties() -
getConfigProperties
@SchemaFieldDescription("Properties passed to the Hadoop configuration the catalog uses.") public abstract @Nullable Map<String,String> getConfigProperties() -
getTriggeringFrequencySeconds
@SchemaFieldDescription("For a streaming pipeline, sets the frequency at which incoming files are appended (default 600, or 10min).") public abstract @Nullable Integer getTriggeringFrequencySeconds() -
getManifestFileSize
@SchemaFieldDescription("The number of data files per manifest (default 10,000 files).") public abstract @Nullable Integer getManifestFileSize() -
getLocationPrefix
@SchemaFieldDescription("The prefix shared among all partitions. For example, a data file may have the following location:%n\'gs://bucket/namespace/table/data/id=13/name=beam/data_file.parquet\'%n%nThe provided prefix should go up until the partition information:%n\'gs://bucket/namespace/table/data/\'.%nIf not provided, will try determining each DataFile\'s partition from its metrics metadata.") public abstract @Nullable String getLocationPrefix() -
getPartitionFields
@SchemaFieldDescription("Fields used to create a partition spec that is applied when tables are created. For a field \'foo\', the available partition transforms are:\n\n- `foo`\n- `truncate(foo, N)`\n- `bucket(foo, N)`\n- `hour(foo)`\n- `day(foo)`\n- `month(foo)`\n- `year(foo)`\n- `void(foo)`\n\nFor more information on partition transforms, please visit https://iceberg.apache.org/spec/#partition-transforms.") public abstract @Nullable List<String> getPartitionFields() -
getTableProperties
@SchemaFieldDescription("Iceberg table properties to be set on the table when it is created.\nFor more information on table properties, please visit https://iceberg.apache.org/docs/latest/configuration/#table-properties.") public abstract @Nullable Map<String,String> getTableProperties() -
getSortFields
@SchemaFieldDescription("Fields used to set the table\'s sort order, applied when the table is created. Each entry has the form `<term> [asc|desc] [nulls first|nulls last]`, where `<term>` is a field name or one of the partition transforms (e.g. `bucket(col, 4)`, `day(ts)`). Direction defaults to ascending; null order defaults to nulls-first for ascending and nulls-last for descending.\nFor more information on sort orders, please visit https://iceberg.apache.org/spec/#sort-orders.") public abstract @Nullable List<String> getSortFields() -
getSchemaEvolutionOptions
@SchemaFieldDescription("Lets the transform change the table schema so that the table has a column for every column the files have. Values: ALLOW_FIELD_ADDITION (columns a file has and the table lacks are added, as optional), ALLOW_FIELD_RELAXATION (a required table column becomes optional when a file lacks it or may hold nulls in it), ALLOW_TYPE_PROMOTION (a column type is widened, for example int to long). Leave it empty to never change the table schema. When any option is set, the transform reads the footer of every Parquet file and commits the allowed changes before registering any file, so every registered file has statistics for all of its columns. A file that needs a change that is not allowed is incompatible; see incompatible_schema_handling. Only Parquet files can be checked: ORC and Avro files are sent to the error output unless unverifiable_file_handling is ACCEPT. Files sent to the error output are dropped unless error_handling is set. If the table does not exist it is created from the union of the Parquet schemas; with none, every file goes to the error output. Batch pipelines only; streaming pipelines cannot use schema evolution yet.") public abstract @Nullable List<String> getSchemaEvolutionOptions() -
getRequiredColumns
@SchemaFieldDescription("Columns that must always be present and never null, as dotted paths for nested fields (for example address.city). They are never made optional, whatever the options allow, and are created as required when the transform creates the table. A file that lacks one of these columns, or holds nulls in it, is sent to the error output (see error_handling). So is a file whose footer marks the column as optional and has no null-count statistics for it, unless unverifiable_file_handling is ACCEPT. Requires schema_evolution_options.") public abstract @Nullable List<String> getRequiredColumns() -
getDryRun
@SchemaFieldDescription("When true, nothing is committed or registered: the transform reads the files\' schemas and emits a `dry_run_report` output with one row that describes what a real run would do. Its `allowed` field is true when every file schema can be merged and the configuration raises no problem; otherwise its `reason` field says what a real run would do about it (fail, or route the files to the error output). Its `schemas` field lists each distinct file schema with the changes a real run would make for it and, when it cannot be merged, why. The output only exists when this is set; consume it as input: `<this transform\'s name>.dry_run_report`. Against a missing table, a REST catalog needs table-create permission even though no table is created.") public abstract @Nullable Boolean getDryRun() -
getIncompatibleSchemaHandling
@SchemaFieldDescription("What happens when a file\'s schema cannot be made to fit the table: it needs a change that is not allowed, or it conflicts with the table or with another file. FAIL_PIPELINE (the default) fails the pipeline before any schema change is committed. ROUTE_TO_ERRORS commits the changes for the other files and sends the incompatible files to the error output; it requires error_handling.") public abstract @Nullable String getIncompatibleSchemaHandling() -
getUnverifiableFileHandling
@SchemaFieldDescription("What happens to a file the checks cannot verify: an ORC or Avro file (the checks read Parquet footers only), or a Parquet file with no null-count statistics for a required column (statistics disabled by the writer, or a column under a list or map). REJECT (the default) sends the file to the error output (see error_handling). ACCEPT registers it without the checks; such files are counted and logged. A file that fails a check is always sent to the error output. An accepted file that lacks a required column, or holds nulls in it, makes reads of the table fail.") public abstract @Nullable String getUnverifiableFileHandling() -
getErrorHandling
@SchemaFieldDescription("Whether and where to output the files that could not be registered, as rows with the file path and the error. Without it those files are dropped.") public abstract @Nullable ErrorHandling getErrorHandling() -
getSchemaEvolution
Validates and converts the schema evolution settings; null when none are set. -
getIcebergCatalog
-