From 92ccc2d47b8197f539516fd65c03d12e75324dbc Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Mon, 10 Aug 2026 15:56:42 -0700 Subject: [PATCH 1/7] ORC/Spark: forward-port default-value reads via idToConstant (LI #76) Restore the production Spark ORC initial-default path from f20062316: omit absent defaulted fields from the ORC projection and inject them into idToConstant so row and vectorized readers fill via existing constant readers. Adapted to the upstream initial-default API and gated behind supportsInitialDefaults so shared ORCSchemaUtil does not break Generic. --- .../main/java/org/apache/iceberg/orc/ORC.java | 17 +- .../org/apache/iceberg/orc/ORCSchemaUtil.java | 43 +++- .../org/apache/iceberg/orc/OrcIterable.java | 9 +- .../iceberg/orc/OrcSchemaWithTypeVisitor.java | 9 +- .../spark/OrcSchemaWithTypeVisitorSpark.java | 127 ++++++++++ .../iceberg/spark/data/SparkOrcReader.java | 8 +- .../vectorized/ConstantArrayColumnVector.java | 136 +++++++++++ .../data/vectorized/ConstantColumnVector.java | 41 +++- .../vectorized/VectorizedSparkOrcReaders.java | 8 +- .../iceberg/spark/source/BaseDataReader.java | 76 ++++-- .../iceberg/spark/source/BatchDataReader.java | 1 + .../iceberg/spark/source/RowDataReader.java | 1 + ...arkOrcReaderForFieldsWithDefaultValue.java | 226 ++++++++++++++++++ 13 files changed, 661 insertions(+), 41 deletions(-) create mode 100644 spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java create mode 100644 spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java create mode 100644 spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java diff --git a/orc/src/main/java/org/apache/iceberg/orc/ORC.java b/orc/src/main/java/org/apache/iceberg/orc/ORC.java index 89cd1ad436..33e9efffac 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/ORC.java +++ b/orc/src/main/java/org/apache/iceberg/orc/ORC.java @@ -672,6 +672,7 @@ public static class ReadBuilder { private Function> readerFunc; private Function> batchedReaderFunc; + private boolean supportsInitialDefaults = false; private int recordsPerBatch = VectorizedRowBatch.DEFAULT_SIZE; private ReadBuilder(InputFile file) { @@ -725,6 +726,19 @@ public ReadBuilder createReaderFunc(Function> r return this; } + /** + * Signals that the configured reader can fill {@code initial-default} values for fields omitted + * from an ORC file. Disabled by default so existing readers retain null-synthesizing projection + * behavior. Spark ORC readers opt in (forward-port of LI #76). + */ + public ReadBuilder supportsInitialDefaults() { + Preconditions.checkState( + this.readerFunc != null || this.batchedReaderFunc != null, + "A reader function must be configured before enabling initial defaults"); + this.supportsInitialDefaults = true; + return this; + } + public ReadBuilder filter(Expression newFilter) { this.filter = newFilter; return this; @@ -733,7 +747,7 @@ public ReadBuilder filter(Expression newFilter) { public ReadBuilder createBatchedReaderFunc( Function> batchReaderFunction) { Preconditions.checkArgument( - this.readerFunc == null, + this.readerFunc == null && !this.supportsInitialDefaults, "Batched reader function cannot be set since the non-batched version is already set"); this.batchedReaderFunc = batchReaderFunction; return this; @@ -759,6 +773,7 @@ public CloseableIterable build() { start, length, readerFunc, + supportsInitialDefaults, caseSensitive, filter, batchedReaderFunc, diff --git a/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java b/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java index fae1a76c37..f013c1179a 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java +++ b/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java @@ -261,18 +261,46 @@ public static Schema convert(TypeDescription orcSchema) { */ public static TypeDescription buildOrcProjection( Schema schema, TypeDescription originalOrcSchema) { + return buildOrcProjection(schema, originalOrcSchema, false); + } + + /** + * Builds the ORC read projection, optionally omitting absent fields that declare an {@code + * initial-default} so a default-aware reader can fill them via {@code idToConstant}. + * + *

Forward-port of LI #76 ({@code f20062316}): when {@code supportsInitialDefaults} is true and + * a field is absent from the file but declares {@code initialDefault()}, it is omitted from the + * projection instead of being synthesized as a null column. + * + *

TODO(PR2): tighten with EMBEDDED vs name-mapped provenance and complete-id checks before + * omitting. + */ + static TypeDescription buildOrcProjection( + Schema schema, TypeDescription originalOrcSchema, boolean supportsInitialDefaults) { final Map icebergToOrc = icebergToOrcMapping("root", originalOrcSchema); - return buildOrcProjection(Integer.MIN_VALUE, schema.asStruct(), true, icebergToOrc); + return buildOrcProjection( + Integer.MIN_VALUE, schema.asStruct(), true, supportsInitialDefaults, icebergToOrc); } private static TypeDescription buildOrcProjection( - Integer fieldId, Type type, boolean isRequired, Map mapping) { + Integer fieldId, + Type type, + boolean isRequired, + boolean supportsInitialDefaults, + Map mapping) { final TypeDescription orcType; switch (type.typeId()) { case STRUCT: orcType = TypeDescription.createStruct(); for (Types.NestedField nestedField : type.asStructType().fields()) { + // Forward-port of LI #76: omit absent defaulted fields so the reader fills via + // idToConstant instead of synthesizing a null column. + if (supportsInitialDefaults + && mapping.get(nestedField.fieldId()) == null + && nestedField.initialDefault() != null) { + continue; + } // Using suffix _r to avoid potential underlying issues in ORC reader // with reused column names between ORC and Iceberg; // e.g. renaming column c -> d and adding new column d @@ -285,6 +313,7 @@ private static TypeDescription buildOrcProjection( nestedField.fieldId(), nestedField.type(), isRequired && nestedField.isRequired(), + supportsInitialDefaults, mapping); orcType.addField(name, childType); } @@ -296,16 +325,22 @@ private static TypeDescription buildOrcProjection( list.elementId(), list.elementType(), isRequired && list.isElementRequired(), + supportsInitialDefaults, mapping); orcType = TypeDescription.createList(elementType); break; case MAP: Types.MapType map = (Types.MapType) type; TypeDescription keyType = - buildOrcProjection(map.keyId(), map.keyType(), isRequired, mapping); + buildOrcProjection( + map.keyId(), map.keyType(), isRequired, supportsInitialDefaults, mapping); TypeDescription valueType = buildOrcProjection( - map.valueId(), map.valueType(), isRequired && map.isValueRequired(), mapping); + map.valueId(), + map.valueType(), + isRequired && map.isValueRequired(), + supportsInitialDefaults, + mapping); orcType = TypeDescription.createMap(keyType, valueType); break; default: diff --git a/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java b/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java index 58cf5d1f96..85a47994a9 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java +++ b/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java @@ -46,6 +46,7 @@ class OrcIterable extends CloseableGroup implements CloseableIterable { private final Long start; private final Long length; private final Function> readerFunction; + private final boolean supportsInitialDefaults; private final Expression filter; private final boolean caseSensitive; private final Function> batchReaderFunction; @@ -60,12 +61,14 @@ class OrcIterable extends CloseableGroup implements CloseableIterable { Long start, Long length, Function> readerFunction, + boolean supportsInitialDefaults, boolean caseSensitive, Expression filter, Function> batchReaderFunction, int recordsPerBatch) { this.schema = schema; this.readerFunction = readerFunction; + this.supportsInitialDefaults = supportsInitialDefaults; this.file = file; this.nameMapping = nameMapping; this.start = start; @@ -86,13 +89,15 @@ public CloseableIterator iterator() { TypeDescription fileSchema = orcFileReader.getSchema(); final TypeDescription readOrcSchema; if (ORCSchemaUtil.hasIds(fileSchema)) { - readOrcSchema = ORCSchemaUtil.buildOrcProjection(schema, fileSchema); + readOrcSchema = + ORCSchemaUtil.buildOrcProjection(schema, fileSchema, supportsInitialDefaults); } else { if (nameMapping == null) { nameMapping = MappingUtil.create(schema); } TypeDescription typeWithIds = ORCSchemaUtil.applyNameMapping(fileSchema, nameMapping); - readOrcSchema = ORCSchemaUtil.buildOrcProjection(schema, typeWithIds); + readOrcSchema = + ORCSchemaUtil.buildOrcProjection(schema, typeWithIds, supportsInitialDefaults); } SearchArgument sarg = null; diff --git a/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java b/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java index fd37283a86..dc587d9679 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java +++ b/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java @@ -36,7 +36,7 @@ public static T visit( Type iType, TypeDescription schema, OrcSchemaWithTypeVisitor visitor) { switch (schema.getCategory()) { case STRUCT: - return visitRecord(iType != null ? iType.asStructType() : null, schema, visitor); + return visitor.visitRecord(iType != null ? iType.asStructType() : null, schema, visitor); case UNION: throw new UnsupportedOperationException("Cannot handle " + schema); @@ -61,7 +61,12 @@ public static T visit( } } - private static T visitRecord( + /** + * Visits a struct. Overridden by Spark to inject {@code initial-default} values into {@code + * idToConstant} for fields omitted from the ORC projection (forward-port of LI #76 / {@code + * f20062316}). + */ + protected T visitRecord( Types.StructType struct, TypeDescription record, OrcSchemaWithTypeVisitor visitor) { List fields = record.getChildren(); List names = record.getFieldNames(); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java new file mode 100644 index 0000000000..da1eb87ee7 --- /dev/null +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java @@ -0,0 +1,127 @@ +/* + * 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.iceberg.spark; + +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.iceberg.MetadataColumns; +import org.apache.iceberg.orc.ORCSchemaUtil; +import org.apache.iceberg.orc.OrcSchemaWithTypeVisitor; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.spark.source.BaseDataReader; +import org.apache.iceberg.types.Types; +import org.apache.orc.TypeDescription; + +/** + * Spark ORC schema visitor that injects {@code initial-default} values into {@code idToConstant} + * for fields omitted from the per-file ORC projection. + * + *

Forward-port of LI #76 ({@code f20062316} / Raymond Zhang). Both {@link + * org.apache.iceberg.spark.data.SparkOrcReader} and {@link + * org.apache.iceberg.spark.data.vectorized.VectorizedSparkOrcReaders} extend this so row and + * vectorized paths share the same inject. + * + *

TODO(PR4): extract inject into {@code ORCSchemaUtil.withInitialDefaults} and delete this + * visitor so Generic/other engines can reuse the same helper without a Spark subclass. + */ +public abstract class OrcSchemaWithTypeVisitorSpark extends OrcSchemaWithTypeVisitor { + + private final Map idToConstant; + + public Map getIdToConstant() { + return idToConstant; + } + + protected OrcSchemaWithTypeVisitorSpark(Map idToConstant) { + this.idToConstant = Maps.newHashMap(); + this.idToConstant.putAll(idToConstant); + } + + @Override + protected T visitRecord( + Types.StructType struct, TypeDescription record, OrcSchemaWithTypeVisitor visitor) { + Preconditions.checkState( + icebergFieldIdsContainOrcFieldIdsInOrder(struct, record), + "Iceberg schema and ORC schema doesn't align, please call ORCSchemaUtil.buildOrcProjection" + + " to get an aligned ORC schema first!"); + List iFields = struct.fields(); + List fields = record.getChildren(); + List names = record.getFieldNames(); + List results = Lists.newArrayListWithExpectedSize(fields.size()); + + for (int i = 0, j = 0; i < iFields.size(); i++) { + Types.NestedField iField = iFields.get(i); + TypeDescription field = j < fields.size() ? fields.get(j) : null; + if (field == null || (iField.fieldId() != ORCSchemaUtil.fieldId(field))) { + // Cases that use idToConstant for an iField: + // 1. MetadataColumns.ROW_POSITION → RowPositionReader + // 2. Partition column → ConstantReader (already in idToConstant from PartitionUtil) + // 3. Field omitted because it declares initial-default → ConstantReader (inject here) + if (MetadataColumns.nonMetadataColumn(iField.name()) + && !idToConstant.containsKey(iField.fieldId()) + && iField.initialDefault() != null) { + idToConstant.put( + iField.fieldId(), + BaseDataReader.convertConstant(iField.type(), iField.initialDefault())); + } + } else { + results.add(visit(iField.type(), field, visitor)); + j++; + } + } + return visitor.record(struct, record, names, results); + } + + private static boolean icebergFieldIdsContainOrcFieldIdsInOrder( + Types.StructType struct, TypeDescription record) { + List icebergIDList = + struct.fields().stream().map(Types.NestedField::fieldId).collect(Collectors.toList()); + List orcIDList = + record.getChildren().stream().map(ORCSchemaUtil::fieldId).collect(Collectors.toList()); + + return containsInOrder(icebergIDList, orcIDList); + } + + /** + * Checks whether {@code list1} contains all integers from {@code list2} in the same relative + * order. {@code list1} may contain extra integers that {@code list2} does not. + */ + private static boolean containsInOrder(List list1, List list2) { + if (list1.size() < list2.size()) { + return false; + } + + for (int i = 0, j = 0; j < list2.size(); j++) { + if (i >= list1.size()) { + return false; + } + while (!list1.get(i).equals(list2.get(j))) { + i++; + if (i >= list1.size()) { + return false; + } + } + i++; + } + return true; + } +} diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/SparkOrcReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/SparkOrcReader.java index 3d53233538..f102213dee 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/SparkOrcReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/SparkOrcReader.java @@ -25,6 +25,7 @@ import org.apache.iceberg.orc.OrcValueReader; import org.apache.iceberg.orc.OrcValueReaders; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.OrcSchemaWithTypeVisitorSpark; import org.apache.iceberg.types.Type; import org.apache.iceberg.types.Types; import org.apache.orc.TypeDescription; @@ -64,11 +65,10 @@ public void setBatchContext(long batchOffsetInFile) { reader.setBatchContext(batchOffsetInFile); } - private static class ReadBuilder extends OrcSchemaWithTypeVisitor> { - private final Map idToConstant; + private static class ReadBuilder extends OrcSchemaWithTypeVisitorSpark> { private ReadBuilder(Map idToConstant) { - this.idToConstant = idToConstant; + super(idToConstant); } @Override @@ -77,7 +77,7 @@ public OrcValueReader record( TypeDescription record, List names, List> fields) { - return SparkOrcValueReaders.struct(record, fields, expected, idToConstant); + return SparkOrcValueReaders.struct(record, fields, expected, getIdToConstant()); } @Override diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java new file mode 100644 index 0000000000..87376c6c8f --- /dev/null +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java @@ -0,0 +1,136 @@ +/* + * 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.iceberg.spark.data.vectorized; + +import java.util.Arrays; +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.catalyst.util.ArrayData; +import org.apache.spark.sql.catalyst.util.MapData; +import org.apache.spark.sql.types.ArrayType; +import org.apache.spark.sql.types.DataType; +import org.apache.spark.sql.types.Decimal; +import org.apache.spark.sql.types.MapType; +import org.apache.spark.sql.types.StructType; +import org.apache.spark.sql.vectorized.ColumnVector; +import org.apache.spark.sql.vectorized.ColumnarArray; +import org.apache.spark.sql.vectorized.ColumnarMap; +import org.apache.spark.unsafe.types.UTF8String; + +/** + * Constant array column vector from LI #76 ({@code f20062316}). + * + *

Unused until nested-typed defaults are legal in the API (TODO PR7/PR8). Kept so the + * forward-port does not delete Raymond's nested vectorized path. + */ +public class ConstantArrayColumnVector extends ConstantColumnVector { + + private final Object[] constantArray; + + public ConstantArrayColumnVector(DataType type, int batchSize, Object[] constantArray) { + super(type, batchSize, constantArray); + this.constantArray = constantArray; + } + + @Override + public boolean getBoolean(int rowId) { + return (boolean) constantArray[rowId]; + } + + @Override + public byte getByte(int rowId) { + return (byte) constantArray[rowId]; + } + + @Override + public short getShort(int rowId) { + return (short) constantArray[rowId]; + } + + @Override + public int getInt(int rowId) { + return (int) constantArray[rowId]; + } + + @Override + public long getLong(int rowId) { + return (long) constantArray[rowId]; + } + + @Override + public float getFloat(int rowId) { + return (float) constantArray[rowId]; + } + + @Override + public double getDouble(int rowId) { + return (double) constantArray[rowId]; + } + + @Override + public Decimal getDecimal(int rowId, int precision, int scale) { + return (Decimal) constantArray[rowId]; + } + + @Override + public UTF8String getUTF8String(int rowId) { + return (UTF8String) constantArray[rowId]; + } + + @Override + public byte[] getBinary(int rowId) { + return (byte[]) constantArray[rowId]; + } + + @Override + public ColumnarArray getArray(int rowId) { + return new ColumnarArray( + new ConstantArrayColumnVector( + ((ArrayType) type).elementType(), + getBatchSize(), + ((ArrayData) constantArray[rowId]).array()), + 0, + ((ArrayData) constantArray[rowId]).numElements()); + } + + @Override + public ColumnarMap getMap(int rowId) { + ColumnVector keys = + new ConstantArrayColumnVector( + ((MapType) type).keyType(), + getBatchSize(), + ((MapData) constantArray[rowId]).keyArray().array()); + ColumnVector values = + new ConstantArrayColumnVector( + ((MapType) type).valueType(), + getBatchSize(), + ((MapData) constantArray[rowId]).valueArray().array()); + return new ColumnarMap(keys, values, 0, ((MapData) constantArray[rowId]).numElements()); + } + + @Override + public ColumnVector getChild(int ordinal) { + DataType fieldType = ((StructType) type).fields()[ordinal].dataType(); + return new ConstantArrayColumnVector( + fieldType, + getBatchSize(), + Arrays.stream(constantArray) + .map(row -> ((InternalRow) row).get(ordinal, fieldType)) + .toArray()); + } +} diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java index 42683ffa90..d4487ab3bb 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java @@ -20,12 +20,25 @@ import org.apache.iceberg.spark.SparkSchemaUtil; import org.apache.iceberg.types.Type; +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.catalyst.util.ArrayData; +import org.apache.spark.sql.catalyst.util.MapData; +import org.apache.spark.sql.types.ArrayType; +import org.apache.spark.sql.types.DataType; import org.apache.spark.sql.types.Decimal; +import org.apache.spark.sql.types.MapType; +import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.vectorized.ColumnVector; import org.apache.spark.sql.vectorized.ColumnarArray; import org.apache.spark.sql.vectorized.ColumnarMap; import org.apache.spark.unsafe.types.UTF8String; +/** + * Constant column vector for partition values and initial-defaults. + * + *

Nested getArray/getMap/getChild support is retained from LI #76 for nested-typed defaults + * (TODO PR7/PR8: re-enable once {@code castDefault} allows non-null nested defaults). + */ class ConstantColumnVector extends ColumnVector { private final Object constant; @@ -37,6 +50,16 @@ class ConstantColumnVector extends ColumnVector { this.batchSize = batchSize; } + ConstantColumnVector(DataType type, int batchSize, Object constant) { + super(type); + this.constant = constant; + this.batchSize = batchSize; + } + + protected int getBatchSize() { + return batchSize; + } + @Override public void close() {} @@ -92,12 +115,22 @@ public double getDouble(int rowId) { @Override public ColumnarArray getArray(int rowId) { - throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); + return new ColumnarArray( + new ConstantArrayColumnVector( + ((ArrayType) type).elementType(), batchSize, ((ArrayData) constant).array()), + 0, + ((ArrayData) constant).numElements()); } @Override public ColumnarMap getMap(int ordinal) { - throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); + ColumnVector keys = + new ConstantArrayColumnVector( + ((MapType) type).keyType(), batchSize, ((MapData) constant).keyArray().array()); + ColumnVector values = + new ConstantArrayColumnVector( + ((MapType) type).valueType(), batchSize, ((MapData) constant).valueArray().array()); + return new ColumnarMap(keys, values, 0, ((MapData) constant).numElements()); } @Override @@ -117,6 +150,8 @@ public byte[] getBinary(int rowId) { @Override public ColumnVector getChild(int ordinal) { - throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); + DataType fieldType = ((StructType) type).fields()[ordinal].dataType(); + return new ConstantColumnVector( + fieldType, batchSize, ((InternalRow) constant).get(ordinal, fieldType)); } } diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java index 7c3b825a62..5a59f84d70 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java @@ -28,6 +28,7 @@ import org.apache.iceberg.orc.OrcValueReader; import org.apache.iceberg.orc.OrcValueReaders; import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.spark.OrcSchemaWithTypeVisitorSpark; import org.apache.iceberg.spark.SparkSchemaUtil; import org.apache.iceberg.spark.data.SparkOrcValueReaders; import org.apache.iceberg.types.Type; @@ -85,11 +86,10 @@ ColumnVector convert( long batchOffsetInFile); } - private static class ReadBuilder extends OrcSchemaWithTypeVisitor { - private final Map idToConstant; + private static class ReadBuilder extends OrcSchemaWithTypeVisitorSpark { private ReadBuilder(Map idToConstant) { - this.idToConstant = idToConstant; + super(idToConstant); } @Override @@ -98,7 +98,7 @@ public Converter record( TypeDescription record, List names, List fields) { - return new StructConverter(iStruct, fields, idToConstant); + return new StructConverter(iStruct, fields, getIdToConstant()); } @Override diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java index f3ddd50eef..14296ea265 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java @@ -22,9 +22,11 @@ import java.io.IOException; import java.math.BigDecimal; import java.nio.ByteBuffer; +import java.util.Collection; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import java.util.stream.Stream; import org.apache.avro.generic.GenericData; import org.apache.avro.util.Utf8; @@ -41,6 +43,7 @@ import org.apache.iceberg.io.InputFile; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.types.Type; import org.apache.iceberg.types.Types.NestedField; @@ -48,9 +51,14 @@ import org.apache.iceberg.util.ByteBuffers; import org.apache.iceberg.util.PartitionUtil; import org.apache.spark.rdd.InputFileBlockHolder; +import org.apache.spark.sql.catalyst.InternalRow; import org.apache.spark.sql.catalyst.expressions.GenericInternalRow; +import org.apache.spark.sql.catalyst.util.ArrayBasedMapData; +import org.apache.spark.sql.catalyst.util.ArrayData; +import org.apache.spark.sql.catalyst.util.GenericArrayData; import org.apache.spark.sql.types.Decimal; import org.apache.spark.unsafe.types.UTF8String; +import scala.collection.JavaConverters; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -59,7 +67,7 @@ * * @param is the Java class returned by this reader whose objects contain one or more rows. */ -abstract class BaseDataReader implements Closeable { +public abstract class BaseDataReader implements Closeable { private static final Logger LOG = LoggerFactory.getLogger(BaseDataReader.class); private final Table table; @@ -160,12 +168,57 @@ protected InputFile getInputFile(String location) { } } - protected static Object convertConstant(Type type, Object value) { + /** + * Converts a constant (partition value or initial-default) to Spark's in-memory representation. + * + *

List/map/struct branches are retained from LI #76 for nested-typed defaults; those defaults + * cannot be declared yet because {@code NestedField.castDefault} rejects non-null nested types + * (TODO PR7/PR8). + */ + public static Object convertConstant(Type type, Object value) { if (value == null) { return null; } switch (type.typeId()) { + case STRUCT: + StructType structType = type.asStructType(); + + if (structType.fields().isEmpty()) { + return new GenericInternalRow(); + } + + InternalRow ret = new GenericInternalRow(structType.fields().size()); + for (int i = 0; i < structType.fields().size(); i++) { + NestedField field = structType.fields().get(i); + Type fieldType = field.type(); + if (value instanceof Map) { + ret.update(i, convertConstant(field.type(), ((Map) value).get(field.name()))); + } else { + ret.update( + i, + convertConstant( + fieldType, ((StructLike) value).get(i, fieldType.typeId().javaClass()))); + } + } + return ret; + case LIST: + List javaList = + ((Collection) value) + .stream() + .map(e -> convertConstant(type.asListType().elementType(), e)) + .collect(Collectors.toList()); + return ArrayData.toArrayData( + JavaConverters.collectionAsScalaIterableConverter(javaList).asScala().toSeq()); + case MAP: + List keyList = Lists.newArrayList(); + List valueList = Lists.newArrayList(); + for (Map.Entry entry : ((Map) value).entrySet()) { + keyList.add(convertConstant(type.asMapType().keyType(), entry.getKey())); + valueList.add(convertConstant(type.asMapType().valueType(), entry.getValue())); + } + return new ArrayBasedMapData( + new GenericArrayData(keyList.toArray()), new GenericArrayData(valueList.toArray())); case DECIMAL: return Decimal.apply((BigDecimal) value); case STRING: @@ -183,25 +236,6 @@ protected static Object convertConstant(Type type, Object value) { return ByteBuffers.toByteArray((ByteBuffer) value); case BINARY: return ByteBuffers.toByteArray((ByteBuffer) value); - case STRUCT: - StructType structType = (StructType) type; - - if (structType.fields().isEmpty()) { - return new GenericInternalRow(); - } - - List fields = structType.fields(); - Object[] values = new Object[fields.size()]; - StructLike struct = (StructLike) value; - - for (int index = 0; index < fields.size(); index++) { - NestedField field = fields.get(index); - Type fieldType = field.type(); - values[index] = - convertConstant(fieldType, struct.get(index, fieldType.typeId().javaClass())); - } - - return new GenericInternalRow(values); default: } return value; diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java index 68e98ba913..dffdd8bf3b 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java @@ -122,6 +122,7 @@ CloseableIterator open(FileScanTask task) { fileSchema -> VectorizedSparkOrcReaders.buildReader( expectedSchema, fileSchema, idToConstant)) + .supportsInitialDefaults() .recordsPerBatch(batchSize) .filter(task.residual()) .caseSensitive(caseSensitive); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/RowDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/RowDataReader.java index f206149da3..041be30f6e 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/RowDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/RowDataReader.java @@ -159,6 +159,7 @@ private CloseableIterable newOrcIterable( .split(task.start(), task.length()) .createReaderFunc( readOrcSchema -> new SparkOrcReader(readSchema, readOrcSchema, idToConstant)) + .supportsInitialDefaults() .filter(task.residual()) .caseSensitive(caseSensitive); diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java new file mode 100644 index 0000000000..7a522c99a1 --- /dev/null +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java @@ -0,0 +1,226 @@ +/* + * 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.iceberg.spark.data; + +import static org.apache.iceberg.spark.data.TestHelpers.assertEquals; + +import java.io.File; +import java.io.IOException; +import java.util.Iterator; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.iceberg.Files; +import org.apache.iceberg.Schema; +import org.apache.iceberg.expressions.Expressions; +import org.apache.iceberg.io.CloseableIterable; +import org.apache.iceberg.orc.ORC; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.relocated.com.google.common.collect.Iterators; +import org.apache.iceberg.spark.data.vectorized.VectorizedSparkOrcReaders; +import org.apache.iceberg.types.Types; +import org.apache.orc.OrcFile; +import org.apache.orc.TypeDescription; +import org.apache.orc.Writer; +import org.apache.orc.storage.ql.exec.vector.LongColumnVector; +import org.apache.orc.storage.ql.exec.vector.VectorizedRowBatch; +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.catalyst.expressions.GenericInternalRow; +import org.apache.spark.sql.vectorized.ColumnarBatch; +import org.apache.spark.unsafe.types.UTF8String; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TemporaryFolder; + +/** + * Forward-port of LI #76 ({@code f20062316}) Spark ORC default-value reads, adapted to the upstream + * {@code initial-default} API. + * + *

Nested-typed column defaults (list/map/struct values) from Raymond's original test are deferred + * until PR7 (API lift of {@code castDefault}) and PR8 (re-enable ConstantArray path). Filter-on- + * defaulted-column coverage is deferred to PR3 (SARG). + */ +public class TestSparkOrcReaderForFieldsWithDefaultValue { + + @Rule public TemporaryFolder temp = new TemporaryFolder(); + + @Test + public void testOrcScalarDefaultValuesRowAndVectorized() throws IOException { + final int numRows = 10; + + final InternalRow expectedFirstRow = new GenericInternalRow(2); + expectedFirstRow.update(0, 0); + expectedFirstRow.update(1, UTF8String.fromString("foo")); + + TypeDescription orcSchema = TypeDescription.fromString("struct"); + + Schema readSchema = + new Schema( + Types.NestedField.required(1, "col1", Types.IntegerType.get()), + Types.NestedField.optional("col2") + .withId(2) + .ofType(Types.StringType.get()) + .withInitialDefault(Expressions.lit("foo")) + .build()); + + File orcFile = writeOrcWithIntColumn(orcSchema, numRows); + + try (CloseableIterable reader = + ORC.read(Files.localInput(orcFile)) + .project(readSchema) + .createReaderFunc(readOrcSchema -> new SparkOrcReader(readSchema, readOrcSchema)) + .supportsInitialDefaults() + .build()) { + InternalRow actualFirstRow = reader.iterator().next(); + assertEquals(readSchema, expectedFirstRow, actualFirstRow); + } + + try (CloseableIterable reader = + ORC.read(Files.localInput(orcFile)) + .project(readSchema) + .createBatchedReaderFunc( + readOrcSchema -> + VectorizedSparkOrcReaders.buildReader( + readSchema, readOrcSchema, ImmutableMap.of())) + .supportsInitialDefaults() + .build()) { + InternalRow actualFirstRow = batchesToRows(reader.iterator()).next(); + assertEquals(readSchema, expectedFirstRow, actualFirstRow); + } + } + + @Test + public void testOrcNestedScalarDefaultValuesRowAndVectorized() throws IOException { + // Parent struct `loc` is present in the file; child `country` is new with an initial-default. + // (If the parent itself is absent, the synthetic null parent short-circuits child constants — + // same as Raymond/Option B: nested scalars assume the ancestor struct column exists.) + final int numRows = 10; + + final InternalRow expectedLoc = new GenericInternalRow(1); + expectedLoc.update(0, UTF8String.fromString("US")); + final InternalRow expectedFirstRow = new GenericInternalRow(2); + expectedFirstRow.update(0, 0L); + expectedFirstRow.update(1, expectedLoc); + + // Empty loc struct in the file: country is absent and will be filled from initial-default. + TypeDescription orcSchema = TypeDescription.fromString("struct>"); + + Schema readSchema = + new Schema( + Types.NestedField.required(1, "id", Types.LongType.get()), + Types.NestedField.optional("loc") + .withId(2) + .ofType( + Types.StructType.of( + Types.NestedField.optional("country") + .withId(3) + .ofType(Types.StringType.get()) + .withInitialDefault(Expressions.lit("US")) + .build())) + .build()); + + File orcFile = writeOrcWithIdAndEmptyLoc(orcSchema, numRows); + + try (CloseableIterable reader = + ORC.read(Files.localInput(orcFile)) + .project(readSchema) + .createReaderFunc(readOrcSchema -> new SparkOrcReader(readSchema, readOrcSchema)) + .supportsInitialDefaults() + .build()) { + InternalRow actualFirstRow = reader.iterator().next(); + assertEquals(readSchema, expectedFirstRow, actualFirstRow); + } + + try (CloseableIterable reader = + ORC.read(Files.localInput(orcFile)) + .project(readSchema) + .createBatchedReaderFunc( + readOrcSchema -> + VectorizedSparkOrcReaders.buildReader( + readSchema, readOrcSchema, ImmutableMap.of())) + .supportsInitialDefaults() + .build()) { + InternalRow actualFirstRow = batchesToRows(reader.iterator()).next(); + assertEquals(readSchema, expectedFirstRow, actualFirstRow); + } + } + + private File writeOrcWithIntColumn(TypeDescription orcSchema, int numRows) throws IOException { + Configuration conf = new Configuration(); + File orcFile = temp.newFile(); + Path orcFilePath = new Path(orcFile.getPath()); + + Writer writer = + OrcFile.createWriter( + orcFilePath, OrcFile.writerOptions(conf).setSchema(orcSchema).overwrite(true)); + + VectorizedRowBatch batch = orcSchema.createRowBatch(); + LongColumnVector firstCol = (LongColumnVector) batch.cols[0]; + for (int r = 0; r < numRows; ++r) { + int row = batch.size++; + firstCol.vector[row] = r; + if (batch.size == batch.getMaxSize()) { + writer.addRowBatch(batch); + batch.reset(); + } + } + if (batch.size != 0) { + writer.addRowBatch(batch); + batch.reset(); + } + writer.close(); + return orcFile; + } + + private File writeOrcWithIdAndEmptyLoc(TypeDescription orcSchema, int numRows) + throws IOException { + Configuration conf = new Configuration(); + File orcFile = temp.newFile(); + Path orcFilePath = new Path(orcFile.getPath()); + + Writer writer = + OrcFile.createWriter( + orcFilePath, OrcFile.writerOptions(conf).setSchema(orcSchema).overwrite(true)); + + VectorizedRowBatch batch = orcSchema.createRowBatch(); + LongColumnVector idCol = (LongColumnVector) batch.cols[0]; + org.apache.orc.storage.ql.exec.vector.StructColumnVector locCol = + (org.apache.orc.storage.ql.exec.vector.StructColumnVector) batch.cols[1]; + for (int r = 0; r < numRows; ++r) { + int row = batch.size++; + idCol.vector[row] = r; + // non-null empty loc struct + locCol.noNulls = true; + locCol.isNull[row] = false; + if (batch.size == batch.getMaxSize()) { + writer.addRowBatch(batch); + batch.reset(); + } + } + if (batch.size != 0) { + writer.addRowBatch(batch); + batch.reset(); + } + writer.close(); + return orcFile; + } + + private Iterator batchesToRows(Iterator batches) { + return Iterators.concat(Iterators.transform(batches, ColumnarBatch::rowIterator)); + } +} From 9c90f1a1de290d7becee826c5ea565f505e72454 Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Mon, 10 Aug 2026 16:07:38 -0700 Subject: [PATCH 2/7] ORC/Spark: apply spotless formatting for CI --- orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java | 3 +-- .../org/apache/iceberg/spark/source/BaseDataReader.java | 2 +- .../data/TestSparkOrcReaderForFieldsWithDefaultValue.java | 6 +++--- 3 files changed, 5 insertions(+), 6 deletions(-) diff --git a/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java b/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java index 85a47994a9..cd8148fdb5 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java +++ b/orc/src/main/java/org/apache/iceberg/orc/OrcIterable.java @@ -89,8 +89,7 @@ public CloseableIterator iterator() { TypeDescription fileSchema = orcFileReader.getSchema(); final TypeDescription readOrcSchema; if (ORCSchemaUtil.hasIds(fileSchema)) { - readOrcSchema = - ORCSchemaUtil.buildOrcProjection(schema, fileSchema, supportsInitialDefaults); + readOrcSchema = ORCSchemaUtil.buildOrcProjection(schema, fileSchema, supportsInitialDefaults); } else { if (nameMapping == null) { nameMapping = MappingUtil.create(schema); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java index 14296ea265..17637ba42e 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java @@ -58,9 +58,9 @@ import org.apache.spark.sql.catalyst.util.GenericArrayData; import org.apache.spark.sql.types.Decimal; import org.apache.spark.unsafe.types.UTF8String; -import scala.collection.JavaConverters; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import scala.collection.JavaConverters; /** * Base class of Spark readers. diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java index 7a522c99a1..9270290acb 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java @@ -51,9 +51,9 @@ * Forward-port of LI #76 ({@code f20062316}) Spark ORC default-value reads, adapted to the upstream * {@code initial-default} API. * - *

Nested-typed column defaults (list/map/struct values) from Raymond's original test are deferred - * until PR7 (API lift of {@code castDefault}) and PR8 (re-enable ConstantArray path). Filter-on- - * defaulted-column coverage is deferred to PR3 (SARG). + *

Nested-typed column defaults (list/map/struct values) from Raymond's original test are + * deferred until PR7 (API lift of {@code castDefault}) and PR8 (re-enable ConstantArray path). + * Filter-on- defaulted-column coverage is deferred to PR3 (SARG). */ public class TestSparkOrcReaderForFieldsWithDefaultValue { From 006163bf6f2bce8b5e12ce0a4e1c81ca552a7f1a Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Wed, 12 Aug 2026 12:15:58 -0700 Subject: [PATCH 3/7] ORC/Spark: keep defaulted projections on the row reader Vectorized ORC default-fill is out of scope. SparkBatchScan disables columnar ORC when the projection includes an initial-default, and the batched reader no longer opts into omission. --- .../spark/OrcSchemaWithTypeVisitorSpark.java | 8 +-- .../iceberg/spark/source/BatchDataReader.java | 1 - .../iceberg/spark/source/SparkBatchScan.java | 12 +++- ...arkOrcReaderForFieldsWithDefaultValue.java | 42 ++--------- .../TestSparkBatchScanInitialDefaults.java | 69 +++++++++++++++++++ 5 files changed, 88 insertions(+), 44 deletions(-) create mode 100644 spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkBatchScanInitialDefaults.java diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java index da1eb87ee7..a89f998fc2 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java @@ -35,10 +35,10 @@ * Spark ORC schema visitor that injects {@code initial-default} values into {@code idToConstant} * for fields omitted from the per-file ORC projection. * - *

Forward-port of LI #76 ({@code f20062316} / Raymond Zhang). Both {@link - * org.apache.iceberg.spark.data.SparkOrcReader} and {@link - * org.apache.iceberg.spark.data.vectorized.VectorizedSparkOrcReaders} extend this so row and - * vectorized paths share the same inject. + *

Forward-port of LI #76 ({@code f20062316} / Raymond Zhang). {@link + * org.apache.iceberg.spark.data.SparkOrcReader} uses this inject. Vectorized ORC default-fill is + * out of scope: Spark routes defaulted projections to the row reader, and the batched ORC reader + * does not opt into omission. * *

TODO(PR4): extract inject into {@code ORCSchemaUtil.withInitialDefaults} and delete this * visitor so Generic/other engines can reuse the same helper without a Spark subclass. diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java index dffdd8bf3b..68e98ba913 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BatchDataReader.java @@ -122,7 +122,6 @@ CloseableIterator open(FileScanTask task) { fileSchema -> VectorizedSparkOrcReaders.buildReader( expectedSchema, fileSchema, idToConstant)) - .supportsInitialDefaults() .recordsPerBatch(batchSize) .filter(task.residual()) .caseSensitive(caseSensitive); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java index 63489d5056..675542f972 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java @@ -38,6 +38,7 @@ import org.apache.iceberg.spark.SparkReadConf; import org.apache.iceberg.spark.SparkSchemaUtil; import org.apache.iceberg.spark.SparkUtil; +import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.TableScanUtil; import org.apache.iceberg.util.Tasks; @@ -193,7 +194,11 @@ public PartitionReaderFactory createReaderFactory() { boolean batchReadsEnabled = batchReadsEnabled(allParquetFileScanTasks, allOrcFileScanTasks); - boolean batchReadOrc = hasNoDeleteFiles && allOrcFileScanTasks; + boolean hasNoInitialDefaults = hasNoInitialDefaults(expectedSchema); + + // Vectorized ORC default-fill is out of scope for this stack. Defaulted projections stay on + // the row reader; the batched ORC path is unchanged for scans that do not project a default. + boolean batchReadOrc = hasNoDeleteFiles && allOrcFileScanTasks && hasNoInitialDefaults; boolean batchReadParquet = hasNoEqDeleteFiles && allParquetFileScanTasks && atLeastOneColumn && onlyPrimitives; @@ -205,6 +210,11 @@ public PartitionReaderFactory createReaderFactory() { return new ReaderFactory(batchSize); } + static boolean hasNoInitialDefaults(Schema schema) { + return TypeUtil.indexById(schema.asStruct()).values().stream() + .noneMatch(field -> field.initialDefault() != null); + } + private boolean batchReadsEnabled(boolean isParquetOnly, boolean isOrcOnly) { if (isParquetOnly) { return readConf.parquetVectorizationEnabled(); diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java index 9270290acb..ef1f9dc487 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java @@ -22,7 +22,6 @@ import java.io.File; import java.io.IOException; -import java.util.Iterator; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.iceberg.Files; @@ -30,9 +29,6 @@ import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.orc.ORC; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.relocated.com.google.common.collect.Iterators; -import org.apache.iceberg.spark.data.vectorized.VectorizedSparkOrcReaders; import org.apache.iceberg.types.Types; import org.apache.orc.OrcFile; import org.apache.orc.TypeDescription; @@ -41,7 +37,6 @@ import org.apache.orc.storage.ql.exec.vector.VectorizedRowBatch; import org.apache.spark.sql.catalyst.InternalRow; import org.apache.spark.sql.catalyst.expressions.GenericInternalRow; -import org.apache.spark.sql.vectorized.ColumnarBatch; import org.apache.spark.unsafe.types.UTF8String; import org.junit.Rule; import org.junit.Test; @@ -49,7 +44,8 @@ /** * Forward-port of LI #76 ({@code f20062316}) Spark ORC default-value reads, adapted to the upstream - * {@code initial-default} API. + * {@code initial-default} API. Spark scans that project a default stay on the row reader; + * vectorized ORC default-fill is out of scope. * *

Nested-typed column defaults (list/map/struct values) from Raymond's original test are * deferred until PR7 (API lift of {@code castDefault}) and PR8 (re-enable ConstantArray path). @@ -60,7 +56,7 @@ public class TestSparkOrcReaderForFieldsWithDefaultValue { @Rule public TemporaryFolder temp = new TemporaryFolder(); @Test - public void testOrcScalarDefaultValuesRowAndVectorized() throws IOException { + public void testOrcScalarDefaultValues() throws IOException { final int numRows = 10; final InternalRow expectedFirstRow = new GenericInternalRow(2); @@ -89,23 +85,10 @@ public void testOrcScalarDefaultValuesRowAndVectorized() throws IOException { InternalRow actualFirstRow = reader.iterator().next(); assertEquals(readSchema, expectedFirstRow, actualFirstRow); } - - try (CloseableIterable reader = - ORC.read(Files.localInput(orcFile)) - .project(readSchema) - .createBatchedReaderFunc( - readOrcSchema -> - VectorizedSparkOrcReaders.buildReader( - readSchema, readOrcSchema, ImmutableMap.of())) - .supportsInitialDefaults() - .build()) { - InternalRow actualFirstRow = batchesToRows(reader.iterator()).next(); - assertEquals(readSchema, expectedFirstRow, actualFirstRow); - } } @Test - public void testOrcNestedScalarDefaultValuesRowAndVectorized() throws IOException { + public void testOrcNestedScalarDefaultValues() throws IOException { // Parent struct `loc` is present in the file; child `country` is new with an initial-default. // (If the parent itself is absent, the synthetic null parent short-circuits child constants — // same as Raymond/Option B: nested scalars assume the ancestor struct column exists.) @@ -145,19 +128,6 @@ public void testOrcNestedScalarDefaultValuesRowAndVectorized() throws IOExceptio InternalRow actualFirstRow = reader.iterator().next(); assertEquals(readSchema, expectedFirstRow, actualFirstRow); } - - try (CloseableIterable reader = - ORC.read(Files.localInput(orcFile)) - .project(readSchema) - .createBatchedReaderFunc( - readOrcSchema -> - VectorizedSparkOrcReaders.buildReader( - readSchema, readOrcSchema, ImmutableMap.of())) - .supportsInitialDefaults() - .build()) { - InternalRow actualFirstRow = batchesToRows(reader.iterator()).next(); - assertEquals(readSchema, expectedFirstRow, actualFirstRow); - } } private File writeOrcWithIntColumn(TypeDescription orcSchema, int numRows) throws IOException { @@ -219,8 +189,4 @@ private File writeOrcWithIdAndEmptyLoc(TypeDescription orcSchema, int numRows) writer.close(); return orcFile; } - - private Iterator batchesToRows(Iterator batches) { - return Iterators.concat(Iterators.transform(batches, ColumnarBatch::rowIterator)); - } } diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkBatchScanInitialDefaults.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkBatchScanInitialDefaults.java new file mode 100644 index 0000000000..7d0fc98074 --- /dev/null +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkBatchScanInitialDefaults.java @@ -0,0 +1,69 @@ +/* + * 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.iceberg.spark.source; + +import org.apache.iceberg.Schema; +import org.apache.iceberg.expressions.Expressions; +import org.apache.iceberg.types.Types; +import org.junit.Assert; +import org.junit.Test; + +public class TestSparkBatchScanInitialDefaults { + + @Test + public void testDisablesVectorizationWhenTopLevelDefaultIsProjected() { + Schema schema = + new Schema( + Types.NestedField.required(1, "id", Types.LongType.get()), + Types.NestedField.optional("country") + .withId(2) + .ofType(Types.StringType.get()) + .withInitialDefault(Expressions.lit("US")) + .build()); + + Assert.assertFalse(SparkBatchScan.hasNoInitialDefaults(schema)); + } + + @Test + public void testDisablesVectorizationWhenNestedDefaultIsProjected() { + Schema schema = + new Schema( + Types.NestedField.required( + 1, + "location", + Types.StructType.of( + Types.NestedField.optional("country") + .withId(2) + .ofType(Types.StringType.get()) + .withInitialDefault(Expressions.lit("US")) + .build()))); + + Assert.assertFalse(SparkBatchScan.hasNoInitialDefaults(schema)); + } + + @Test + public void testAllowsVectorizationWhenNoDefaultIsProjected() { + Schema schema = + new Schema( + Types.NestedField.required(1, "id", Types.LongType.get()), + Types.NestedField.optional(2, "country", Types.StringType.get())); + + Assert.assertTrue(SparkBatchScan.hasNoInitialDefaults(schema)); + } +} From 7d2b5f08edb4166c40bc02ba3b3d55eb5be4d713 Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Wed, 12 Aug 2026 13:28:18 -0700 Subject: [PATCH 4/7] ORC/Spark: trim process notes from source comments Keep the why (opt-in omit, row-only fill, nested defaults not API-legal). Drop stack TODOs, commit SHAs, and forward-port narration. --- orc/src/main/java/org/apache/iceberg/orc/ORC.java | 2 +- .../java/org/apache/iceberg/orc/ORCSchemaUtil.java | 11 +++-------- .../apache/iceberg/orc/OrcSchemaWithTypeVisitor.java | 3 +-- .../iceberg/spark/OrcSchemaWithTypeVisitorSpark.java | 9 ++------- .../data/vectorized/ConstantArrayColumnVector.java | 7 +------ .../spark/data/vectorized/ConstantColumnVector.java | 4 ++-- .../apache/iceberg/spark/source/BaseDataReader.java | 5 ++--- .../apache/iceberg/spark/source/SparkBatchScan.java | 3 +-- .../TestSparkOrcReaderForFieldsWithDefaultValue.java | 12 +++--------- 9 files changed, 16 insertions(+), 40 deletions(-) diff --git a/orc/src/main/java/org/apache/iceberg/orc/ORC.java b/orc/src/main/java/org/apache/iceberg/orc/ORC.java index 33e9efffac..bfe00e788e 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/ORC.java +++ b/orc/src/main/java/org/apache/iceberg/orc/ORC.java @@ -729,7 +729,7 @@ public ReadBuilder createReaderFunc(Function> r /** * Signals that the configured reader can fill {@code initial-default} values for fields omitted * from an ORC file. Disabled by default so existing readers retain null-synthesizing projection - * behavior. Spark ORC readers opt in (forward-port of LI #76). + * behavior. */ public ReadBuilder supportsInitialDefaults() { Preconditions.checkState( diff --git a/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java b/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java index f013c1179a..e641c53dd4 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java +++ b/orc/src/main/java/org/apache/iceberg/orc/ORCSchemaUtil.java @@ -268,12 +268,8 @@ public static TypeDescription buildOrcProjection( * Builds the ORC read projection, optionally omitting absent fields that declare an {@code * initial-default} so a default-aware reader can fill them via {@code idToConstant}. * - *

Forward-port of LI #76 ({@code f20062316}): when {@code supportsInitialDefaults} is true and - * a field is absent from the file but declares {@code initialDefault()}, it is omitted from the - * projection instead of being synthesized as a null column. - * - *

TODO(PR2): tighten with EMBEDDED vs name-mapped provenance and complete-id checks before - * omitting. + *

When {@code supportsInitialDefaults} is true and a field is absent from the file but + * declares {@code initialDefault()}, it is omitted instead of being synthesized as a null column. */ static TypeDescription buildOrcProjection( Schema schema, TypeDescription originalOrcSchema, boolean supportsInitialDefaults) { @@ -294,8 +290,7 @@ private static TypeDescription buildOrcProjection( case STRUCT: orcType = TypeDescription.createStruct(); for (Types.NestedField nestedField : type.asStructType().fields()) { - // Forward-port of LI #76: omit absent defaulted fields so the reader fills via - // idToConstant instead of synthesizing a null column. + // Omit so the reader fills via idToConstant instead of a synthetic null column. if (supportsInitialDefaults && mapping.get(nestedField.fieldId()) == null && nestedField.initialDefault() != null) { diff --git a/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java b/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java index dc587d9679..568de02198 100644 --- a/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java +++ b/orc/src/main/java/org/apache/iceberg/orc/OrcSchemaWithTypeVisitor.java @@ -63,8 +63,7 @@ public static T visit( /** * Visits a struct. Overridden by Spark to inject {@code initial-default} values into {@code - * idToConstant} for fields omitted from the ORC projection (forward-port of LI #76 / {@code - * f20062316}). + * idToConstant} for fields omitted from the ORC projection. */ protected T visitRecord( Types.StructType struct, TypeDescription record, OrcSchemaWithTypeVisitor visitor) { diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java index a89f998fc2..c14dfef1aa 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java @@ -35,13 +35,8 @@ * Spark ORC schema visitor that injects {@code initial-default} values into {@code idToConstant} * for fields omitted from the per-file ORC projection. * - *

Forward-port of LI #76 ({@code f20062316} / Raymond Zhang). {@link - * org.apache.iceberg.spark.data.SparkOrcReader} uses this inject. Vectorized ORC default-fill is - * out of scope: Spark routes defaulted projections to the row reader, and the batched ORC reader - * does not opt into omission. - * - *

TODO(PR4): extract inject into {@code ORCSchemaUtil.withInitialDefaults} and delete this - * visitor so Generic/other engines can reuse the same helper without a Spark subclass. + *

{@link org.apache.iceberg.spark.data.SparkOrcReader} uses this inject. The batched ORC reader + * does not opt into omission, so defaulted Spark scans stay on the row reader. */ public abstract class OrcSchemaWithTypeVisitorSpark extends OrcSchemaWithTypeVisitor { diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java index 87376c6c8f..192a0015d2 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java @@ -32,12 +32,7 @@ import org.apache.spark.sql.vectorized.ColumnarMap; import org.apache.spark.unsafe.types.UTF8String; -/** - * Constant array column vector from LI #76 ({@code f20062316}). - * - *

Unused until nested-typed defaults are legal in the API (TODO PR7/PR8). Kept so the - * forward-port does not delete Raymond's nested vectorized path. - */ +/** Constant array column vector. Unused until nested-typed defaults are legal in the API. */ public class ConstantArrayColumnVector extends ConstantColumnVector { private final Object[] constantArray; diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java index d4487ab3bb..e8b1a7a8fe 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java @@ -36,8 +36,8 @@ /** * Constant column vector for partition values and initial-defaults. * - *

Nested getArray/getMap/getChild support is retained from LI #76 for nested-typed defaults - * (TODO PR7/PR8: re-enable once {@code castDefault} allows non-null nested defaults). + *

Nested getArray/getMap/getChild support is unused until {@code castDefault} allows non-null + * nested defaults. */ class ConstantColumnVector extends ColumnVector { diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java index 17637ba42e..1b70db6c14 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/BaseDataReader.java @@ -171,9 +171,8 @@ protected InputFile getInputFile(String location) { /** * Converts a constant (partition value or initial-default) to Spark's in-memory representation. * - *

List/map/struct branches are retained from LI #76 for nested-typed defaults; those defaults - * cannot be declared yet because {@code NestedField.castDefault} rejects non-null nested types - * (TODO PR7/PR8). + *

List/map/struct branches are unused until {@code NestedField.castDefault} allows non-null + * nested types. */ public static Object convertConstant(Type type, Object value) { if (value == null) { diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java index 675542f972..6788424781 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkBatchScan.java @@ -196,8 +196,7 @@ public PartitionReaderFactory createReaderFactory() { boolean hasNoInitialDefaults = hasNoInitialDefaults(expectedSchema); - // Vectorized ORC default-fill is out of scope for this stack. Defaulted projections stay on - // the row reader; the batched ORC path is unchanged for scans that do not project a default. + // Defaulted projections stay on the row reader; batched ORC does not fill initial-defaults. boolean batchReadOrc = hasNoDeleteFiles && allOrcFileScanTasks && hasNoInitialDefaults; boolean batchReadParquet = diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java index ef1f9dc487..ffd2268af3 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java @@ -43,13 +43,8 @@ import org.junit.rules.TemporaryFolder; /** - * Forward-port of LI #76 ({@code f20062316}) Spark ORC default-value reads, adapted to the upstream - * {@code initial-default} API. Spark scans that project a default stay on the row reader; - * vectorized ORC default-fill is out of scope. - * - *

Nested-typed column defaults (list/map/struct values) from Raymond's original test are - * deferred until PR7 (API lift of {@code castDefault}) and PR8 (re-enable ConstantArray path). - * Filter-on- defaulted-column coverage is deferred to PR3 (SARG). + * Spark ORC reads of fields with {@code initial-default}. Defaulted projections use the row + * reader. Nested-typed column defaults are not covered: {@code castDefault} rejects them. */ public class TestSparkOrcReaderForFieldsWithDefaultValue { @@ -90,8 +85,7 @@ public void testOrcScalarDefaultValues() throws IOException { @Test public void testOrcNestedScalarDefaultValues() throws IOException { // Parent struct `loc` is present in the file; child `country` is new with an initial-default. - // (If the parent itself is absent, the synthetic null parent short-circuits child constants — - // same as Raymond/Option B: nested scalars assume the ancestor struct column exists.) + // If the parent itself is absent, the synthetic null parent short-circuits child constants. final int numRows = 10; final InternalRow expectedLoc = new GenericInternalRow(1); From 95fbc2789799650f908ad605292a259c4e9b00f7 Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Wed, 12 Aug 2026 16:06:20 -0700 Subject: [PATCH 5/7] ORC/Spark: fix spotless wrap in default-value test javadoc --- .../data/TestSparkOrcReaderForFieldsWithDefaultValue.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java index ffd2268af3..0f6de50518 100644 --- a/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java +++ b/spark/v3.1/spark/src/test/java/org/apache/iceberg/spark/data/TestSparkOrcReaderForFieldsWithDefaultValue.java @@ -43,8 +43,8 @@ import org.junit.rules.TemporaryFolder; /** - * Spark ORC reads of fields with {@code initial-default}. Defaulted projections use the row - * reader. Nested-typed column defaults are not covered: {@code castDefault} rejects them. + * Spark ORC reads of fields with {@code initial-default}. Defaulted projections use the row reader. + * Nested-typed column defaults are not covered: {@code castDefault} rejects them. */ public class TestSparkOrcReaderForFieldsWithDefaultValue { From fd10372f850ed61c14ff8687db4d5593080d2fa2 Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Thu, 13 Aug 2026 21:31:08 -0700 Subject: [PATCH 6/7] ORC/Spark: drop unused vectorized default-fill leftovers Vectorized fill is out of scope; SparkBatchScan already keeps defaulted projections on the row reader. Restore the batched visitor and constant vectors to the pre-#76 path. --- .../vectorized/ConstantArrayColumnVector.java | 131 ------------------ .../data/vectorized/ConstantColumnVector.java | 41 +----- .../vectorized/VectorizedSparkOrcReaders.java | 8 +- 3 files changed, 7 insertions(+), 173 deletions(-) delete mode 100644 spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java deleted file mode 100644 index 192a0015d2..0000000000 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantArrayColumnVector.java +++ /dev/null @@ -1,131 +0,0 @@ -/* - * 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.iceberg.spark.data.vectorized; - -import java.util.Arrays; -import org.apache.spark.sql.catalyst.InternalRow; -import org.apache.spark.sql.catalyst.util.ArrayData; -import org.apache.spark.sql.catalyst.util.MapData; -import org.apache.spark.sql.types.ArrayType; -import org.apache.spark.sql.types.DataType; -import org.apache.spark.sql.types.Decimal; -import org.apache.spark.sql.types.MapType; -import org.apache.spark.sql.types.StructType; -import org.apache.spark.sql.vectorized.ColumnVector; -import org.apache.spark.sql.vectorized.ColumnarArray; -import org.apache.spark.sql.vectorized.ColumnarMap; -import org.apache.spark.unsafe.types.UTF8String; - -/** Constant array column vector. Unused until nested-typed defaults are legal in the API. */ -public class ConstantArrayColumnVector extends ConstantColumnVector { - - private final Object[] constantArray; - - public ConstantArrayColumnVector(DataType type, int batchSize, Object[] constantArray) { - super(type, batchSize, constantArray); - this.constantArray = constantArray; - } - - @Override - public boolean getBoolean(int rowId) { - return (boolean) constantArray[rowId]; - } - - @Override - public byte getByte(int rowId) { - return (byte) constantArray[rowId]; - } - - @Override - public short getShort(int rowId) { - return (short) constantArray[rowId]; - } - - @Override - public int getInt(int rowId) { - return (int) constantArray[rowId]; - } - - @Override - public long getLong(int rowId) { - return (long) constantArray[rowId]; - } - - @Override - public float getFloat(int rowId) { - return (float) constantArray[rowId]; - } - - @Override - public double getDouble(int rowId) { - return (double) constantArray[rowId]; - } - - @Override - public Decimal getDecimal(int rowId, int precision, int scale) { - return (Decimal) constantArray[rowId]; - } - - @Override - public UTF8String getUTF8String(int rowId) { - return (UTF8String) constantArray[rowId]; - } - - @Override - public byte[] getBinary(int rowId) { - return (byte[]) constantArray[rowId]; - } - - @Override - public ColumnarArray getArray(int rowId) { - return new ColumnarArray( - new ConstantArrayColumnVector( - ((ArrayType) type).elementType(), - getBatchSize(), - ((ArrayData) constantArray[rowId]).array()), - 0, - ((ArrayData) constantArray[rowId]).numElements()); - } - - @Override - public ColumnarMap getMap(int rowId) { - ColumnVector keys = - new ConstantArrayColumnVector( - ((MapType) type).keyType(), - getBatchSize(), - ((MapData) constantArray[rowId]).keyArray().array()); - ColumnVector values = - new ConstantArrayColumnVector( - ((MapType) type).valueType(), - getBatchSize(), - ((MapData) constantArray[rowId]).valueArray().array()); - return new ColumnarMap(keys, values, 0, ((MapData) constantArray[rowId]).numElements()); - } - - @Override - public ColumnVector getChild(int ordinal) { - DataType fieldType = ((StructType) type).fields()[ordinal].dataType(); - return new ConstantArrayColumnVector( - fieldType, - getBatchSize(), - Arrays.stream(constantArray) - .map(row -> ((InternalRow) row).get(ordinal, fieldType)) - .toArray()); - } -} diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java index e8b1a7a8fe..42683ffa90 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/ConstantColumnVector.java @@ -20,25 +20,12 @@ import org.apache.iceberg.spark.SparkSchemaUtil; import org.apache.iceberg.types.Type; -import org.apache.spark.sql.catalyst.InternalRow; -import org.apache.spark.sql.catalyst.util.ArrayData; -import org.apache.spark.sql.catalyst.util.MapData; -import org.apache.spark.sql.types.ArrayType; -import org.apache.spark.sql.types.DataType; import org.apache.spark.sql.types.Decimal; -import org.apache.spark.sql.types.MapType; -import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.vectorized.ColumnVector; import org.apache.spark.sql.vectorized.ColumnarArray; import org.apache.spark.sql.vectorized.ColumnarMap; import org.apache.spark.unsafe.types.UTF8String; -/** - * Constant column vector for partition values and initial-defaults. - * - *

Nested getArray/getMap/getChild support is unused until {@code castDefault} allows non-null - * nested defaults. - */ class ConstantColumnVector extends ColumnVector { private final Object constant; @@ -50,16 +37,6 @@ class ConstantColumnVector extends ColumnVector { this.batchSize = batchSize; } - ConstantColumnVector(DataType type, int batchSize, Object constant) { - super(type); - this.constant = constant; - this.batchSize = batchSize; - } - - protected int getBatchSize() { - return batchSize; - } - @Override public void close() {} @@ -115,22 +92,12 @@ public double getDouble(int rowId) { @Override public ColumnarArray getArray(int rowId) { - return new ColumnarArray( - new ConstantArrayColumnVector( - ((ArrayType) type).elementType(), batchSize, ((ArrayData) constant).array()), - 0, - ((ArrayData) constant).numElements()); + throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); } @Override public ColumnarMap getMap(int ordinal) { - ColumnVector keys = - new ConstantArrayColumnVector( - ((MapType) type).keyType(), batchSize, ((MapData) constant).keyArray().array()); - ColumnVector values = - new ConstantArrayColumnVector( - ((MapType) type).valueType(), batchSize, ((MapData) constant).valueArray().array()); - return new ColumnarMap(keys, values, 0, ((MapData) constant).numElements()); + throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); } @Override @@ -150,8 +117,6 @@ public byte[] getBinary(int rowId) { @Override public ColumnVector getChild(int ordinal) { - DataType fieldType = ((StructType) type).fields()[ordinal].dataType(); - return new ConstantColumnVector( - fieldType, batchSize, ((InternalRow) constant).get(ordinal, fieldType)); + throw new UnsupportedOperationException("ConstantColumnVector only supports primitives"); } } diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java index 5a59f84d70..7c3b825a62 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/data/vectorized/VectorizedSparkOrcReaders.java @@ -28,7 +28,6 @@ import org.apache.iceberg.orc.OrcValueReader; import org.apache.iceberg.orc.OrcValueReaders; import org.apache.iceberg.relocated.com.google.common.collect.Lists; -import org.apache.iceberg.spark.OrcSchemaWithTypeVisitorSpark; import org.apache.iceberg.spark.SparkSchemaUtil; import org.apache.iceberg.spark.data.SparkOrcValueReaders; import org.apache.iceberg.types.Type; @@ -86,10 +85,11 @@ ColumnVector convert( long batchOffsetInFile); } - private static class ReadBuilder extends OrcSchemaWithTypeVisitorSpark { + private static class ReadBuilder extends OrcSchemaWithTypeVisitor { + private final Map idToConstant; private ReadBuilder(Map idToConstant) { - super(idToConstant); + this.idToConstant = idToConstant; } @Override @@ -98,7 +98,7 @@ public Converter record( TypeDescription record, List names, List fields) { - return new StructConverter(iStruct, fields, getIdToConstant()); + return new StructConverter(iStruct, fields, idToConstant); } @Override From 609e8b989a40dd25610331bfe59d2d703308a553 Mon Sep 17 00:00:00 2001 From: Christian Bush Date: Thu, 13 Aug 2026 21:31:16 -0700 Subject: [PATCH 7/7] ORC/Spark: Spark inject visitor is row-reader only --- .../apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java index c14dfef1aa..95e8e2894a 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/OrcSchemaWithTypeVisitorSpark.java @@ -35,8 +35,8 @@ * Spark ORC schema visitor that injects {@code initial-default} values into {@code idToConstant} * for fields omitted from the per-file ORC projection. * - *

{@link org.apache.iceberg.spark.data.SparkOrcReader} uses this inject. The batched ORC reader - * does not opt into omission, so defaulted Spark scans stay on the row reader. + *

{@link org.apache.iceberg.spark.data.SparkOrcReader} uses this inject. Vectorized ORC does + * not; {@code SparkBatchScan} keeps defaulted projections on the row reader. */ public abstract class OrcSchemaWithTypeVisitorSpark extends OrcSchemaWithTypeVisitor {