This is an automated email from the ASF dual-hosted git repository.

claudevdm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new b1068bc0acf AddFiles: SchemaEvolutionConfig, the public settings for 
the schema pre-pass (#40003)
b1068bc0acf is described below

commit b1068bc0acfcca88c6c1e7331d789b59ec6b8bf0
Author: claudevdm <[email protected]>
AuthorDate: Tue Sep 8 11:53:03 2026 -0400

    AddFiles: SchemaEvolutionConfig, the public settings for the schema 
pre-pass (#40003)
    
    The user-facing shape of schema evolution for AddFiles, added ahead of
    the code that acts on it so the API can be reviewed on its own.
    
    SchemaEvolutionOption is three independent switches: ALLOW_FIELD_ADDITION
    (new columns, added optional), ALLOW_FIELD_RELAXATION (a required column
    becomes optional when a file may hold nulls in it or lacks it), and
    ALLOW_TYPE_PROMOTION (widening along Iceberg's rules). They constrain
    columns the table already has; a column new in a window takes the union
    of that window's file schemas.
    
    SchemaEvolutionConfig is a Serializable AutoValue. With no options it is
    disabled and files register as before this series. Pins
    (getRequiredColumns) are dotted column paths that are never relaxed and
    are created required; a Parquet file that cannot prove they hold no
    nulls goes to the error output. IncompatibleSchemaHandling picks
    FAIL_PIPELINE (batch only, the batch default) or ROUTE_TO_ERRORS (the
    streaming default) for a file schema the options cannot cover.
---
 .../beam/sdk/io/iceberg/SchemaEvolutionConfig.java | 152 +++++++++++++++++++++
 .../beam/sdk/io/iceberg/SchemaEvolutionOption.java |  51 +++++++
 .../sdk/io/iceberg/SchemaEvolutionConfigTest.java  |  92 +++++++++++++
 3 files changed, 295 insertions(+)

diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfig.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfig.java
new file mode 100644
index 00000000000..492b1cdad49
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfig.java
@@ -0,0 +1,152 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import com.google.auto.value.AutoValue;
+import java.io.Serializable;
+import java.util.Collections;
+import java.util.EnumSet;
+import java.util.Set;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Schema evolution settings for {@link 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.
+ *
+ * <pre>{@code
+ * 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();
+ * }</pre>
+ *
+ * <p><b>Pins.</b> 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, has nulls in it, or carries no null-count statistics for it is 
routed to the error
+ * output; ORC and Avro files are not checked. Pins name canonical (table) 
paths, dotted for nested
+ * fields, with the container segment spelled out under lists and maps ({@code
+ * addresses.element.city}, {@code attributes.value.total}). A top-level 
column whose own name
+ * contains a dot cannot be pinned.
+ *
+ * <p><b>Incompatible schemas.</b> A schema that needs a change the options do 
not allow, or that
+ * conflicts with the table or with another file's schema. {@link 
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.
+ */
+@AutoValue
+public abstract class SchemaEvolutionConfig implements Serializable {
+
+  public enum IncompatibleSchemaHandling {
+    /**
+     * Fail the pipeline before committing any schema change, with a message 
listing every
+     * incompatible schema, its reason and file count. The batch default, and 
batch only: in
+     * streaming a failing window's commit would be retried forever and hold 
every later window, so
+     * {@link AddFiles} rejects this setting for unbounded input.
+     */
+    FAIL_PIPELINE,
+    /**
+     * Skip the incompatible schema, commit the rest, and route its files to 
the error output with
+     * the specific reason. The streaming default.
+     */
+    ROUTE_TO_ERRORS
+  }
+
+  public abstract Set<SchemaEvolutionOption> getOptions();
+
+  /**
+   * Canonical column paths (dotted for nested fields) that are never relaxed 
and are created
+   * required; files that cannot prove they hold no nulls in them go to the 
error output.
+   */
+  public abstract Set<String> getRequiredColumns();
+
+  public boolean isPinned(String columnPath) {
+    return getRequiredColumns().contains(columnPath);
+  }
+
+  /**
+   * Unset resolves by mode: {@code FAIL_PIPELINE} in batch, {@code 
ROUTE_TO_ERRORS} in streaming.
+   */
+  public abstract @Nullable IncompatibleSchemaHandling 
getIncompatibleSchemaHandling();
+
+  public IncompatibleSchemaHandling incompatibleSchemaHandling(boolean 
bounded) {
+    IncompatibleSchemaHandling handling = getIncompatibleSchemaHandling();
+    if (handling != null) {
+      return handling;
+    }
+    return bounded
+        ? IncompatibleSchemaHandling.FAIL_PIPELINE
+        : IncompatibleSchemaHandling.ROUTE_TO_ERRORS;
+  }
+
+  public boolean isEnabled() {
+    return !getOptions().isEmpty();
+  }
+
+  public boolean allows(SchemaEvolutionOption option) {
+    return getOptions().contains(option);
+  }
+
+  public static SchemaEvolutionConfig disabled() {
+    return builder().build();
+  }
+
+  public static SchemaEvolutionConfig of(SchemaEvolutionOption... options) {
+    Set<SchemaEvolutionOption> set = 
EnumSet.noneOf(SchemaEvolutionOption.class);
+    Collections.addAll(set, options);
+    return builder().setOptions(set).build();
+  }
+
+  public static Builder builder() {
+    return new AutoValue_SchemaEvolutionConfig.Builder()
+        .setOptions(Collections.emptySet())
+        .setRequiredColumns(Collections.emptySet());
+  }
+
+  @AutoValue.Builder
+  public abstract static class Builder {
+    public abstract Builder setOptions(Set<SchemaEvolutionOption> options);
+
+    public abstract Builder setRequiredColumns(Set<String> requiredColumns);
+
+    public abstract Builder setIncompatibleSchemaHandling(
+        @Nullable IncompatibleSchemaHandling handling);
+
+    abstract SchemaEvolutionConfig autoBuild();
+
+    /** Pins and handling without an option would silently do nothing, so they 
are rejected. */
+    public SchemaEvolutionConfig build() {
+      SchemaEvolutionConfig config = autoBuild();
+      for (String column : config.getRequiredColumns()) {
+        Preconditions.checkArgument(
+            !column.trim().isEmpty() && column.equals(column.trim()),
+            "required column is blank or has surrounding whitespace: '%s'",
+            column);
+      }
+      Preconditions.checkArgument(
+          config.isEnabled()
+              || (config.getRequiredColumns().isEmpty()
+                  && config.getIncompatibleSchemaHandling() == null),
+          "required columns and incompatible schema handling need at least one 
schema evolution"
+              + " option");
+      return config;
+    }
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionOption.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionOption.java
new file mode 100644
index 00000000000..62cf54db725
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionOption.java
@@ -0,0 +1,51 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+/**
+ * Kinds of table schema change {@link AddFiles} may make so that a file's 
columns are covered by
+ * the table schema. Evolution is all or nothing per file schema: a file 
schema needing a change
+ * that is not allowed is incompatible; see {@link 
SchemaEvolutionConfig.IncompatibleSchemaHandling}
+ * for what happens to its files. Pinned columns ({@link
+ * SchemaEvolutionConfig#getRequiredColumns()}) are never relaxed whatever the 
options.
+ *
+ * <p>The options constrain changes to columns the table already has when a 
window's schema commit
+ * starts (the whole input, in batch). A column that is new in that window 
takes the union of the
+ * window's file schemas: its type is the widest among them and it is optional 
unless pinned, so two
+ * files that disagree about a new column never conflict with each other, only 
with the table.
+ */
+public enum SchemaEvolutionOption {
+  /**
+   * Add columns present in files but absent from the table, as optional 
columns at every level (a
+   * pinned column is created required). Also lets this transform create a 
missing table from the
+   * union of the file schemas.
+   */
+  ALLOW_FIELD_ADDITION,
+  /**
+   * Make a required table column optional when a file has nulls in it, 
carries no null-count
+   * statistics for it, or lacks it entirely (every row would read null). For 
a column outside lists
+   * and maps, a file whose footer proves zero nulls never triggers this, 
however its writer
+   * declared the column; under a list or map the declaration is taken as is. 
Relaxation is
+   * table-wide and permanent; pin the columns that must stay required.
+   */
+  ALLOW_FIELD_RELAXATION,
+  /**
+   * Widen a column's type (int to long, float to double, decimal precision) 
when a file needs it.
+   */
+  ALLOW_TYPE_PROMOTION
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfigTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfigTest.java
new file mode 100644
index 00000000000..edc8b338bfc
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/SchemaEvolutionConfigTest.java
@@ -0,0 +1,92 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.util.Collections;
+import java.util.EnumSet;
+import java.util.HashSet;
+import 
org.apache.beam.sdk.io.iceberg.SchemaEvolutionConfig.IncompatibleSchemaHandling;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+@RunWith(JUnit4.class)
+public class SchemaEvolutionConfigTest {
+
+  @Test
+  public void testEnabledOnlyWithOptions() {
+    assertFalse(SchemaEvolutionConfig.disabled().isEnabled());
+
+    SchemaEvolutionConfig config =
+        SchemaEvolutionConfig.of(SchemaEvolutionOption.ALLOW_FIELD_ADDITION);
+
+    assertTrue(config.isEnabled());
+    assertTrue(config.allows(SchemaEvolutionOption.ALLOW_FIELD_ADDITION));
+    assertFalse(config.allows(SchemaEvolutionOption.ALLOW_FIELD_RELAXATION));
+  }
+
+  @Test
+  public void testBlankOrPaddedPinIsRejected() {
+    for (String bad : new String[] {"  ", " id", "address.city\t"}) {
+      SchemaEvolutionConfig.Builder builder =
+          SchemaEvolutionConfig.builder()
+              
.setOptions(EnumSet.of(SchemaEvolutionOption.ALLOW_FIELD_ADDITION))
+              .setRequiredColumns(new 
HashSet<>(Collections.singletonList(bad)));
+
+      IllegalArgumentException e = 
assertThrows(IllegalArgumentException.class, builder::build);
+
+      assertTrue(e.getMessage(), e.getMessage().contains("'" + bad + "'"));
+    }
+  }
+
+  @Test
+  public void testPinsAndHandlingRequireAnOption() {
+    SchemaEvolutionConfig.Builder pinsOnly =
+        
SchemaEvolutionConfig.builder().setRequiredColumns(Collections.singleton("id"));
+    SchemaEvolutionConfig.Builder handlingOnly =
+        SchemaEvolutionConfig.builder()
+            
.setIncompatibleSchemaHandling(IncompatibleSchemaHandling.ROUTE_TO_ERRORS);
+
+    IllegalArgumentException e = assertThrows(IllegalArgumentException.class, 
pinsOnly::build);
+    assertThrows(IllegalArgumentException.class, handlingOnly::build);
+
+    assertTrue(e.getMessage(), e.getMessage().contains("at least one schema 
evolution option"));
+  }
+
+  @Test
+  public void testIncompatibleSchemaHandlingDefaultsByMode() {
+    SchemaEvolutionConfig unset =
+        SchemaEvolutionConfig.of(SchemaEvolutionOption.ALLOW_FIELD_ADDITION);
+    SchemaEvolutionConfig forced =
+        SchemaEvolutionConfig.builder()
+            .setOptions(EnumSet.of(SchemaEvolutionOption.ALLOW_FIELD_ADDITION))
+            
.setIncompatibleSchemaHandling(IncompatibleSchemaHandling.ROUTE_TO_ERRORS)
+            .build();
+
+    assertEquals(IncompatibleSchemaHandling.FAIL_PIPELINE, 
unset.incompatibleSchemaHandling(true));
+    assertEquals(
+        IncompatibleSchemaHandling.ROUTE_TO_ERRORS, 
unset.incompatibleSchemaHandling(false));
+    assertEquals(
+        IncompatibleSchemaHandling.ROUTE_TO_ERRORS, 
forced.incompatibleSchemaHandling(true));
+  }
+}

Reply via email to