diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemas.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemas.java
new file mode 100644
index 000000000000..1b81e008f19d
--- /dev/null
+++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemas.java
@@ -0,0 +1,103 @@
+/*
+ * 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 java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.TreeMap;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.coders.CoderRegistry;
+import org.apache.beam.sdk.coders.KvCoder;
+import org.apache.beam.sdk.coders.ListCoder;
+import org.apache.beam.sdk.coders.MapCoder;
+import org.apache.beam.sdk.coders.StringUtf8Coder;
+import org.apache.beam.sdk.coders.VarLongCoder;
+import org.apache.beam.sdk.transforms.Combine;
+import org.apache.beam.sdk.values.KV;
+
+/**
+ * Collects the distinct schemas among canonical file schema JSONs (see {@link FileSchemas}), with
+ * the number of files per schema, most common first (ties broken by JSON). The commit side applies
+ * schemas in this order, so the schema covering the most files wins a conflict.
+ *
+ *
Inputs are compared as strings, so they must already be canonical.
+ */
+class CollectDistinctSchemas
+ extends Combine.CombineFn, List>> {
+
+ @Override
+ public Map createAccumulator() {
+ return new TreeMap<>();
+ }
+
+ @Override
+ public Map addInput(Map accumulator, String schemaJson) {
+ add(accumulator, schemaJson, 1L);
+ return accumulator;
+ }
+
+ @Override
+ public Map mergeAccumulators(Iterable> accumulators) {
+ Map merged = createAccumulator();
+ for (Map accumulator : accumulators) {
+ for (Map.Entry entry : accumulator.entrySet()) {
+ add(merged, entry.getKey(), entry.getValue());
+ }
+ }
+ return merged;
+ }
+
+ @Override
+ public List> extractOutput(Map accumulator) {
+ List> schemas = new ArrayList<>();
+ for (Map.Entry entry : accumulator.entrySet()) {
+ schemas.add(KV.of(entry.getKey(), entry.getValue()));
+ }
+ schemas.sort(
+ (a, b) -> {
+ int byCount = Long.compare(b.getValue(), a.getValue());
+ if (byCount != 0) {
+ return byCount;
+ }
+ return a.getKey().compareTo(b.getKey());
+ });
+ return schemas;
+ }
+
+ @Override
+ public Coder> getAccumulatorCoder(
+ CoderRegistry registry, Coder inputCoder) {
+ return MapCoder.of(StringUtf8Coder.of(), VarLongCoder.of());
+ }
+
+ @Override
+ public Coder>> getDefaultOutputCoder(
+ CoderRegistry registry, Coder inputCoder) {
+ return ListCoder.of(KvCoder.of(StringUtf8Coder.of(), VarLongCoder.of()));
+ }
+
+ private static void add(Map accumulator, String schemaJson, long count) {
+ Long existing = accumulator.get(schemaJson);
+ if (existing == null) {
+ accumulator.put(schemaJson, count);
+ } else {
+ accumulator.put(schemaJson, existing + count);
+ }
+ }
+}
diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileSchemas.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileSchemas.java
new file mode 100644
index 000000000000..592e11e8c767
--- /dev/null
+++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/FileSchemas.java
@@ -0,0 +1,98 @@
+/*
+ * 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 java.util.ArrayList;
+import java.util.List;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.parquet.ParquetSchemaUtil;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.TypeUtil;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+
+/**
+ * Derives the schema a file contributes to schema inference. The canonical form sorts struct fields
+ * by name at every level and renumbers ids in deterministic order, so files that differ only in
+ * column order produce identical JSON. Ids are positional and meaningless: the commit side
+ * reconciles columns by name.
+ */
+final class FileSchemas {
+ private FileSchemas() {}
+
+ static String canonicalJson(ParquetMetadata footer) {
+ Schema converted = ParquetSchemaUtil.convert(footer.getFileMetaData().getSchema());
+ return SchemaParser.toJson(canonical(converted));
+ }
+
+ static Schema canonical(Schema schema) {
+ Type sorted = TypeUtil.visit(schema.asStruct(), new SortFields());
+ int[] nextId = {0};
+ return TypeUtil.assignFreshIds(new Schema(sorted.asStructType().fields()), () -> ++nextId[0]);
+ }
+
+ /**
+ * Rebuilds every struct with its fields sorted by name; every other attribute (optionality, doc,
+ * defaults) is preserved. Iceberg owns the traversal, so nested types this code has never heard
+ * of (variant, and whatever comes next) are visited rather than silently passed through.
+ */
+ private static class SortFields extends TypeUtil.SchemaVisitor {
+ @Override
+ public Type struct(Types.StructType struct, List fieldTypes) {
+ List rebuilt = new ArrayList<>();
+ for (int i = 0; i < struct.fields().size(); i++) {
+ Types.NestedField field = struct.fields().get(i);
+ rebuilt.add(Types.NestedField.from(field).ofType(fieldTypes.get(i)).build());
+ }
+ rebuilt.sort((a, b) -> a.name().compareTo(b.name()));
+ return Types.StructType.of(rebuilt);
+ }
+
+ @Override
+ public Type field(Types.NestedField field, Type fieldType) {
+ return fieldType;
+ }
+
+ @Override
+ public Type list(Types.ListType list, Type elementType) {
+ if (list.isElementOptional()) {
+ return Types.ListType.ofOptional(list.elementId(), elementType);
+ }
+ return Types.ListType.ofRequired(list.elementId(), elementType);
+ }
+
+ @Override
+ public Type map(Types.MapType map, Type keyType, Type valueType) {
+ if (map.isValueOptional()) {
+ return Types.MapType.ofOptional(map.keyId(), map.valueId(), keyType, valueType);
+ }
+ return Types.MapType.ofRequired(map.keyId(), map.valueId(), keyType, valueType);
+ }
+
+ @Override
+ public Type variant(Types.VariantType variant) {
+ return variant;
+ }
+
+ @Override
+ public Type primitive(Type.PrimitiveType primitive) {
+ return primitive;
+ }
+ }
+}
diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchema.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchema.java
new file mode 100644
index 000000000000..571c0b552441
--- /dev/null
+++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchema.java
@@ -0,0 +1,181 @@
+/*
+ * 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.apache.beam.sdk.metrics.Metrics.counter;
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
+
+import java.util.Collections;
+import java.util.concurrent.Callable;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.transforms.windowing.PaneInfo;
+import org.apache.iceberg.FileFormat;
+import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.joda.time.Instant;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Emits the canonical schema (see {@link FileSchemas}) of every readable Parquet file as JSON.
+ * Unreadable or non-Parquet files contribute nothing.
+ */
+class ReadFooterSchema extends DoFn {
+ private static final Logger LOG = LoggerFactory.getLogger(ReadFooterSchema.class);
+
+ static final int DEFAULT_THREAD_POOL_SIZE = 10;
+ static final int DEFAULT_MAX_IN_FLIGHT_TASKS = 100;
+ static final String FILES_READ_COUNTER = "numFilesRead";
+ static final String SCHEMAS_EMITTED_COUNTER = "numSchemasEmitted";
+ static final String FOOTER_READ_ERRORS_COUNTER = "numFooterReadErrors";
+ private static final Counter numFilesRead = counter(ReadFooterSchema.class, FILES_READ_COUNTER);
+ private static final Counter numSchemasEmitted =
+ counter(ReadFooterSchema.class, SCHEMAS_EMITTED_COUNTER);
+ private static final Counter numFooterReadErrors =
+ counter(ReadFooterSchema.class, FOOTER_READ_ERRORS_COUNTER);
+
+ private final int threadPoolSize;
+ private final int maxInFlightTasks;
+ private transient @MonotonicNonNull BoundedAsyncTasks tasks;
+
+ ReadFooterSchema() {
+ this(DEFAULT_THREAD_POOL_SIZE, DEFAULT_MAX_IN_FLIGHT_TASKS);
+ }
+
+ ReadFooterSchema(int threadPoolSize, int maxInFlightTasks) {
+ this.threadPoolSize = threadPoolSize;
+ this.maxInFlightTasks = maxInFlightTasks;
+ }
+
+ /**
+ * {@code schemaJson} is null when the file contributes no schema. Counters are updated when the
+ * result is delivered, on the processing thread: metrics touched from the executor are lost.
+ */
+ private static class ReadResult {
+ final @Nullable String schemaJson;
+ final boolean footerError;
+ final Instant timestamp;
+ final BoundedWindow window;
+ final PaneInfo paneInfo;
+
+ ReadResult(
+ @Nullable String schemaJson,
+ boolean footerError,
+ Instant timestamp,
+ BoundedWindow window,
+ PaneInfo paneInfo) {
+ this.schemaJson = schemaJson;
+ this.footerError = footerError;
+ this.timestamp = timestamp;
+ this.window = window;
+ this.paneInfo = paneInfo;
+ }
+ }
+
+ @Setup
+ public void setup() {
+ tasks = new BoundedAsyncTasks<>(threadPoolSize, maxInFlightTasks);
+ }
+
+ /** Clears anything left behind if the runner reuses this instance after a failed bundle. */
+ @StartBundle
+ public void startBundle() {
+ checkStateNotNull(tasks).cancelAll();
+ }
+
+ @Teardown
+ public void teardown() {
+ if (tasks != null) {
+ tasks.shutdown();
+ }
+ }
+
+ @ProcessElement
+ public void process(
+ @Element String filePath,
+ @Timestamp Instant timestamp,
+ BoundedWindow window,
+ PaneInfo paneInfo,
+ OutputReceiver output)
+ throws Exception {
+ numFilesRead.inc();
+ Callable task = createReadTask(filePath, timestamp, window, paneInfo);
+ checkStateNotNull(tasks).submit(task, result -> outputResult(result, output));
+ }
+
+ @FinishBundle
+ public void finishBundle(FinishBundleContext context) throws Exception {
+ checkStateNotNull(tasks).awaitAll(result -> outputAtFinish(result, context));
+ }
+
+ private static void outputAtFinish(ReadResult result, FinishBundleContext context) {
+ count(result);
+ if (result.schemaJson != null) {
+ context.output(result.schemaJson, result.timestamp, result.window);
+ }
+ }
+
+ private static void outputResult(ReadResult result, OutputReceiver output) {
+ count(result);
+ if (result.schemaJson != null) {
+ output.outputWindowedValue(
+ result.schemaJson,
+ result.timestamp,
+ Collections.singleton(result.window),
+ result.paneInfo);
+ }
+ }
+
+ private static void count(ReadResult result) {
+ if (result.schemaJson != null) {
+ numSchemasEmitted.inc();
+ }
+ if (result.footerError) {
+ numFooterReadErrors.inc();
+ }
+ }
+
+ private static Callable createReadTask(
+ String filePath, Instant timestamp, BoundedWindow window, PaneInfo paneInfo) {
+ return () -> {
+ FileFormat format;
+ try {
+ format = AddFiles.inferFormat(filePath);
+ } catch (AddFiles.UnknownFormatException e) {
+ return new ReadResult(null, false, timestamp, window, paneInfo);
+ }
+ if (!format.equals(FileFormat.PARQUET)) {
+ return new ReadResult(null, false, timestamp, window, paneInfo);
+ }
+ try {
+ ParquetMetadata footer = ParquetFooters.read(filePath);
+ return new ReadResult(
+ FileSchemas.canonicalJson(footer), false, timestamp, window, paneInfo);
+ } catch (Exception e) {
+ LOG.warn(
+ "Could not read the footer of {}; the file will not contribute to schema inference: {}",
+ filePath,
+ AddFiles.errorMessage(e));
+ return new ReadResult(null, true, timestamp, window, paneInfo);
+ }
+ };
+ }
+}
diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemasTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemasTest.java
new file mode 100644
index 000000000000..206e13acd0a6
--- /dev/null
+++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/CollectDistinctSchemasTest.java
@@ -0,0 +1,131 @@
+/*
+ * 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.apache.iceberg.types.Types.NestedField.optional;
+import static org.apache.iceberg.types.Types.NestedField.required;
+import static org.junit.Assert.assertEquals;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.testing.CoderProperties;
+import org.apache.beam.sdk.testing.PAssert;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.transforms.Combine;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.PCollection;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.types.Types;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+@RunWith(JUnit4.class)
+public class CollectDistinctSchemasTest {
+ @Rule public final TestPipeline pipeline = TestPipeline.create();
+
+ private static final String ID_NAME =
+ json(
+ new Schema(
+ required(1, "id", Types.IntegerType.get()),
+ optional(2, "name", Types.StringType.get())));
+ private static final String NAME_ID =
+ json(
+ new Schema(
+ optional(1, "name", Types.StringType.get()),
+ required(2, "id", Types.IntegerType.get())));
+ private static final String ID_LONG_NAME =
+ json(
+ new Schema(
+ required(1, "id", Types.LongType.get()),
+ optional(2, "name", Types.StringType.get())));
+
+ private final CollectDistinctSchemas fn = new CollectDistinctSchemas();
+
+ @Test
+ public void testDedupsIdenticalSchemas() {
+ assertEquals(Arrays.asList(KV.of(ID_NAME, 3L)), combine(ID_NAME, ID_NAME, ID_NAME));
+ }
+
+ /** Inputs are compared as strings; canonicalization is ReadFooterSchema's job. */
+ @Test
+ public void testDifferentStringsAreDistinct() {
+ List> out = combine(ID_NAME, NAME_ID, ID_LONG_NAME);
+ assertEquals(3, out.size());
+ for (KV entry : out) {
+ assertEquals(Long.valueOf(1L), entry.getValue());
+ }
+ }
+
+ @Test
+ public void testMostCommonFirstThenJson() {
+ List> out = combine(NAME_ID, ID_LONG_NAME, ID_NAME, ID_LONG_NAME, NAME_ID);
+ assertEquals(
+ Arrays.asList(KV.of(ID_LONG_NAME, 2L), KV.of(NAME_ID, 2L), KV.of(ID_NAME, 1L)), out);
+ }
+
+ @Test
+ public void testMergeSumsCounts() {
+ Map first = fn.addInput(fn.createAccumulator(), ID_NAME);
+ Map second = fn.addInput(fn.createAccumulator(), ID_NAME);
+ second = fn.addInput(second, NAME_ID);
+ List> out =
+ fn.extractOutput(fn.mergeAccumulators(Arrays.asList(first, second)));
+ assertEquals(Arrays.asList(KV.of(ID_NAME, 2L), KV.of(NAME_ID, 1L)), out);
+ }
+
+ @Test
+ public void testEmptyInput() {
+ assertEquals(Arrays.asList(), combine());
+ }
+
+ @Test
+ public void testAccumulatorCoderRoundTrip() throws Exception {
+ Coder> coder = fn.getAccumulatorCoder(null, null);
+ Map accumulator = fn.addInput(fn.createAccumulator(), ID_NAME);
+ accumulator = fn.addInput(accumulator, NAME_ID);
+ CoderProperties.coderDecodeEncodeEqual(coder, accumulator);
+ }
+
+ @Test
+ public void testPipeline() {
+ PCollection>> out =
+ pipeline
+ .apply(Create.of(ID_NAME, NAME_ID, ID_NAME))
+ .apply(Combine.globally(new CollectDistinctSchemas()));
+ PAssert.that(out).containsInAnyOrder(Arrays.asList(KV.of(ID_NAME, 2L), KV.of(NAME_ID, 1L)));
+ pipeline.run();
+ }
+
+ private List> combine(String... schemaJsons) {
+ Map accumulator = fn.createAccumulator();
+ for (String schemaJson : schemaJsons) {
+ accumulator = fn.addInput(accumulator, schemaJson);
+ }
+ return fn.extractOutput(accumulator);
+ }
+
+ private static String json(Schema schema) {
+ return SchemaParser.toJson(schema);
+ }
+}
diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/FileSchemasTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/FileSchemasTest.java
new file mode 100644
index 000000000000..9723695186be
--- /dev/null
+++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/FileSchemasTest.java
@@ -0,0 +1,131 @@
+/*
+ * 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.apache.iceberg.types.Types.NestedField.optional;
+import static org.apache.iceberg.types.Types.NestedField.required;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.types.Types;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+@RunWith(JUnit4.class)
+public class FileSchemasTest {
+
+ @Test
+ public void testSortsTopLevelFieldsAndRenumbers() {
+ Schema input =
+ new Schema(
+ optional(7, "name", Types.StringType.get()), required(3, "id", Types.LongType.get()));
+ Schema expected =
+ new Schema(
+ required(1, "id", Types.LongType.get()), optional(2, "name", Types.StringType.get()));
+
+ assertSame(expected, FileSchemas.canonical(input));
+ }
+
+ @Test
+ public void testSortsNestedStructFields() {
+ Schema input =
+ new Schema(
+ optional(
+ 1,
+ "address",
+ Types.StructType.of(
+ optional(2, "zip", Types.IntegerType.get()),
+ optional(3, "city", Types.StringType.get()))));
+ Schema expected =
+ new Schema(
+ optional(
+ 1,
+ "address",
+ Types.StructType.of(
+ optional(2, "city", Types.StringType.get()),
+ optional(3, "zip", Types.IntegerType.get()))));
+
+ assertSame(expected, FileSchemas.canonical(input));
+ }
+
+ @Test
+ public void testPermutationsProduceIdenticalJson() {
+ Schema a =
+ new Schema(
+ optional(1, "b", Types.StringType.get()),
+ optional(2, "a", Types.StructType.of(optional(3, "y", Types.LongType.get()))),
+ optional(4, "c", Types.ListType.ofOptional(5, Types.StringType.get())));
+ Schema b =
+ new Schema(
+ optional(1, "c", Types.ListType.ofOptional(2, Types.StringType.get())),
+ optional(3, "a", Types.StructType.of(optional(4, "y", Types.LongType.get()))),
+ optional(5, "b", Types.StringType.get()));
+
+ assertEquals(
+ SchemaParser.toJson(FileSchemas.canonical(a)),
+ SchemaParser.toJson(FileSchemas.canonical(b)));
+ }
+
+ /** Ids number every field of a struct before descending into nested types. */
+ @Test
+ public void testPreservesListMapStructNestingAndOptionality() {
+ Schema input =
+ new Schema(
+ required(
+ 1,
+ "m",
+ Types.MapType.ofRequired(
+ 2,
+ 3,
+ Types.StringType.get(),
+ Types.StructType.of(
+ optional(4, "z", Types.IntegerType.get()),
+ required(5, "a", Types.ListType.ofRequired(6, Types.DoubleType.get()))))));
+ Schema expected =
+ new Schema(
+ required(
+ 1,
+ "m",
+ Types.MapType.ofRequired(
+ 2,
+ 3,
+ Types.StringType.get(),
+ Types.StructType.of(
+ required(4, "a", Types.ListType.ofRequired(6, Types.DoubleType.get())),
+ optional(5, "z", Types.IntegerType.get())))));
+
+ assertSame(expected, FileSchemas.canonical(input));
+ }
+
+ @Test
+ public void testCanonicalSchemaIsUnchanged() {
+ Schema canonical =
+ new Schema(
+ required(1, "a", Types.LongType.get()),
+ optional(2, "b", Types.StructType.of(optional(3, "x", Types.StringType.get()))));
+
+ assertSame(canonical, FileSchemas.canonical(canonical));
+ }
+
+ private static void assertSame(Schema expected, Schema actual) {
+ assertTrue("expected " + expected + " but was " + actual, expected.sameSchema(actual));
+ }
+}
diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchemaTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchemaTest.java
new file mode 100644
index 000000000000..6b3679913918
--- /dev/null
+++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ReadFooterSchemaTest.java
@@ -0,0 +1,279 @@
+/*
+ * 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.apache.iceberg.types.Types.NestedField.optional;
+import static org.apache.iceberg.types.Types.NestedField.required;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import org.apache.avro.SchemaBuilder;
+import org.apache.beam.sdk.PipelineResult;
+import org.apache.beam.sdk.metrics.MetricNameFilter;
+import org.apache.beam.sdk.metrics.MetricQueryResults;
+import org.apache.beam.sdk.metrics.MetricResult;
+import org.apache.beam.sdk.metrics.MetricsFilter;
+import org.apache.beam.sdk.testing.PAssert;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.values.PCollection;
+import org.apache.hadoop.fs.Path;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetWriter;
+import org.apache.iceberg.io.DataWriter;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.avro.AvroParquetWriter;
+import org.apache.parquet.hadoop.ParquetWriter;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+@RunWith(JUnit4.class)
+public class ReadFooterSchemaTest {
+ @Rule public final TestPipeline pipeline = TestPipeline.create();
+ @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+
+ private static final Schema FLAT_SCHEMA =
+ new Schema(
+ required(1, "id", Types.IntegerType.get()), optional(2, "name", Types.StringType.get()));
+
+ private static final Schema NESTED_SCHEMA =
+ new Schema(
+ required(1, "id", Types.LongType.get()),
+ optional(
+ 2,
+ "address",
+ Types.StructType.of(
+ optional(3, "city", Types.StringType.get()),
+ optional(4, "zip", Types.IntegerType.get()))),
+ optional(5, "tags", Types.ListType.ofOptional(6, Types.StringType.get())),
+ optional(
+ 7,
+ "attributes",
+ Types.MapType.ofOptional(8, 9, Types.StringType.get(), Types.DoubleType.get())));
+
+ @Test
+ public void testEmitsFlatSchema() throws IOException {
+ String file = writeParquet("flat.parquet", FLAT_SCHEMA, record(FLAT_SCHEMA, "id", 1));
+
+ assertSchemas(run(file), FLAT_SCHEMA);
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testEmitsNestedSchema() throws IOException {
+ String file = writeParquet("nested.parquet", NESTED_SCHEMA, record(NESTED_SCHEMA, "id", 1L));
+
+ assertSchemas(run(file), NESTED_SCHEMA);
+
+ pipeline.run();
+ }
+
+ /** Iceberg's writer creates no file for zero rows, so this uses parquet-avro directly. */
+ @Test
+ public void testZeroRowParquetEmitsSchema() throws IOException {
+ String file = new File(temporaryFolder.getRoot(), "empty.parquet").getAbsolutePath();
+ org.apache.avro.Schema avroSchema =
+ SchemaBuilder.record("flat").fields().requiredInt("id").optionalString("name").endRecord();
+ ParquetWriter writer =
+ AvroParquetWriter.builder(new Path(file)).withSchema(avroSchema).build();
+ writer.close();
+
+ assertSchemas(run(file), FLAT_SCHEMA);
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testMissingFileEmitsNothing() {
+ String file = new File(temporaryFolder.getRoot(), "missing.parquet").getAbsolutePath();
+
+ PAssert.that(run(file)).empty();
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testGarbageBytesEmitNothing() throws IOException {
+ String file =
+ writeBytes("garbage.parquet", "not a parquet file".getBytes(StandardCharsets.UTF_8));
+
+ PAssert.that(run(file)).empty();
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testTruncatedParquetEmitsNothing() throws IOException {
+ String good = writeParquet("good.parquet", FLAT_SCHEMA, record(FLAT_SCHEMA, "id", 1));
+ byte[] bytes = Files.readAllBytes(new File(good).toPath());
+ String file = writeBytes("truncated.parquet", Arrays.copyOf(bytes, bytes.length / 2));
+
+ PAssert.that(run(file)).empty();
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testZeroByteFileEmitsNothing() throws IOException {
+ String file = writeBytes("zero.parquet", new byte[0]);
+
+ PAssert.that(run(file)).empty();
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testNonParquetEmitsNothing() throws IOException {
+ String avro = writeBytes("data.avro", new byte[0]);
+ String unknown = writeBytes("data.txt", new byte[0]);
+
+ PAssert.that(run(avro, unknown)).empty();
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testPermutedColumnsProduceIdenticalSchema() throws IOException {
+ Schema permuted =
+ new Schema(
+ optional(1, "name", Types.StringType.get()),
+ required(2, "id", Types.IntegerType.get()));
+ String a = writeParquet("a.parquet", FLAT_SCHEMA, record(FLAT_SCHEMA, "id", 1));
+ String b = writeParquet("b.parquet", permuted, record(permuted, "id", 1));
+
+ PAssert.that(run(a, b))
+ .satisfies(
+ actual -> {
+ List jsons = new ArrayList<>();
+ actual.forEach(jsons::add);
+ assertEquals(2, jsons.size());
+ assertEquals(jsons.get(0), jsons.get(1));
+ return null;
+ });
+
+ pipeline.run();
+ }
+
+ @Test
+ public void testMixedBundleEmitsOnlyReadableSchemas() throws IOException {
+ String flat = writeParquet("flat.parquet", FLAT_SCHEMA, record(FLAT_SCHEMA, "id", 1));
+ String nested = writeParquet("nested.parquet", NESTED_SCHEMA, record(NESTED_SCHEMA, "id", 1L));
+ String garbage = writeBytes("garbage.parquet", "garbage".getBytes(StandardCharsets.UTF_8));
+ String avro = writeBytes("data.avro", new byte[0]);
+ String missing = new File(temporaryFolder.getRoot(), "missing.parquet").getAbsolutePath();
+
+ assertSchemas(run(flat, nested, garbage, avro, missing), FLAT_SCHEMA, NESTED_SCHEMA);
+
+ PipelineResult result = pipeline.run();
+
+ assertEquals(5L, counter(result, ReadFooterSchema.FILES_READ_COUNTER));
+ assertEquals(2L, counter(result, ReadFooterSchema.SCHEMAS_EMITTED_COUNTER));
+ assertEquals(2L, counter(result, ReadFooterSchema.FOOTER_READ_ERRORS_COUNTER));
+ }
+
+ private static long counter(PipelineResult result, String name) {
+ MetricQueryResults metrics =
+ result
+ .metrics()
+ .queryMetrics(
+ MetricsFilter.builder()
+ .addNameFilter(MetricNameFilter.named(ReadFooterSchema.class, name))
+ .build());
+ long total = 0;
+ for (MetricResult counter : metrics.getCounters()) {
+ total += counter.getAttempted();
+ }
+ return total;
+ }
+
+ private PCollection run(String... paths) {
+ return pipeline.apply(Create.of(Arrays.asList(paths))).apply(ParDo.of(new ReadFooterSchema()));
+ }
+
+ /** Asserts the emitted schemas equal the canonical forms of {@code expected}, in any order. */
+ private static void assertSchemas(PCollection out, Schema... expected) {
+ List expectedJson = new ArrayList<>();
+ for (Schema schema : expected) {
+ expectedJson.add(SchemaParser.toJson(FileSchemas.canonical(schema)));
+ }
+ PAssert.that(out)
+ .satisfies(
+ actual -> {
+ List remaining = new ArrayList<>(expectedJson);
+ for (String json : actual) {
+ Schema schema = SchemaParser.fromJson(json);
+ boolean matched = false;
+ for (int i = 0; i < remaining.size(); i++) {
+ if (SchemaParser.fromJson(remaining.get(i)).sameSchema(schema)) {
+ remaining.remove(i);
+ matched = true;
+ break;
+ }
+ }
+ assertTrue("Unexpected schema: " + json, matched);
+ }
+ assertEquals("Missing schemas: " + remaining, 0, remaining.size());
+ return null;
+ });
+ }
+
+ private String writeParquet(String name, Schema schema, Record... records) throws IOException {
+ String file = new File(temporaryFolder.getRoot(), name).getAbsolutePath();
+ DataWriter writer =
+ Parquet.writeData(org.apache.iceberg.Files.localOutput(file))
+ .schema(schema)
+ .withSpec(PartitionSpec.unpartitioned())
+ .createWriterFunc(GenericParquetWriter::create)
+ .build();
+ try {
+ for (Record record : records) {
+ writer.write(record);
+ }
+ } finally {
+ writer.close();
+ }
+ return file;
+ }
+
+ private String writeBytes(String name, byte[] bytes) throws IOException {
+ File file = new File(temporaryFolder.getRoot(), name);
+ Files.write(file.toPath(), bytes);
+ return file.getAbsolutePath();
+ }
+
+ private static Record record(Schema schema, String field, Object value) {
+ return GenericRecord.create(schema).copy(field, value);
+ }
+}