-
Notifications
You must be signed in to change notification settings - Fork 45
ORC: bind struct fields by Iceberg field id in generic and Spark 3.1 readers #256
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,220 @@ | ||
| /* | ||
| * 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.data.orc; | ||
|
|
||
| import static org.apache.iceberg.types.Types.NestedField.optional; | ||
| import static org.apache.iceberg.types.Types.NestedField.required; | ||
|
|
||
| import java.io.File; | ||
| import java.io.IOException; | ||
| import java.math.BigDecimal; | ||
| import java.util.List; | ||
| import java.util.Map; | ||
| import org.apache.iceberg.Files; | ||
| import org.apache.iceberg.MetadataColumns; | ||
| import org.apache.iceberg.Schema; | ||
| import org.apache.iceberg.data.GenericRecord; | ||
| import org.apache.iceberg.data.Record; | ||
| import org.apache.iceberg.io.CloseableIterable; | ||
| import org.apache.iceberg.io.FileAppender; | ||
| 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.Lists; | ||
| import org.apache.iceberg.types.TypeUtil; | ||
| import org.apache.iceberg.types.Types; | ||
| import org.junit.Assert; | ||
| import org.junit.Rule; | ||
| import org.junit.Test; | ||
| import org.junit.rules.TemporaryFolder; | ||
|
|
||
| /** | ||
| * Tests specifically targeting the id-based field binding path of {@link GenericOrcReaders}. These | ||
| * exercise branches of the id-binding {@code StructReader} constructor that the shared projection | ||
| * tests do not, and assert that id-binding produces the same results positional binding would. | ||
| */ | ||
| public class TestGenericOrcReaderIdBinding { | ||
|
|
||
| @Rule public TemporaryFolder temp = new TemporaryFolder(); | ||
|
|
||
| private List<Record> writeAndRead( | ||
| String desc, Schema writeSchema, Schema readSchema, List<Record> records) throws IOException { | ||
| return writeAndRead(desc, writeSchema, readSchema, records, ImmutableMap.of()); | ||
| } | ||
|
|
||
| private List<Record> writeAndRead( | ||
| String desc, | ||
| Schema writeSchema, | ||
| Schema readSchema, | ||
| List<Record> records, | ||
| Map<Integer, ?> idToConstant) | ||
| throws IOException { | ||
| File file = temp.newFile(desc + ".orc"); | ||
| Assert.assertTrue("Delete should succeed", file.delete()); | ||
|
|
||
| try (FileAppender<Record> appender = | ||
| ORC.write(Files.localOutput(file)) | ||
| .schema(writeSchema) | ||
| .createWriterFunc(GenericOrcWriter::buildWriter) | ||
| .build()) { | ||
| appender.addAll(records); | ||
| } | ||
|
|
||
| // Project only fields that physically exist in the file (excluding metadata columns) so that | ||
| // buildOrcProjection succeeds, but hand the reader the full read schema so the id-binding | ||
| // constructor resolves metadata/constant fields. | ||
| Schema projection = TypeUtil.selectNot(readSchema, MetadataColumns.metadataFieldIds()); | ||
| try (CloseableIterable<Record> reader = | ||
| ORC.read(Files.localInput(file)) | ||
| .project(projection) | ||
| .createReaderFunc( | ||
| fileSchema -> GenericOrcReader.buildReader(readSchema, fileSchema, idToConstant)) | ||
| .build()) { | ||
| return Lists.newArrayList(reader); | ||
| } | ||
| } | ||
|
|
||
| @Test | ||
| public void testTypePromotion() throws IOException { | ||
| Schema writeSchema = | ||
| new Schema( | ||
| required(1, "id", Types.IntegerType.get()), | ||
| optional(2, "f", Types.FloatType.get()), | ||
| required(3, "dec", Types.DecimalType.of(9, 2))); | ||
|
|
||
| Record record = GenericRecord.create(writeSchema.asStruct()); | ||
| record.setField("id", 42); | ||
| record.setField("f", 1.5f); | ||
| record.setField("dec", new BigDecimal("123.45")); | ||
|
|
||
| Schema promotedSchema = | ||
| new Schema( | ||
| required(1, "id", Types.LongType.get()), | ||
| optional(2, "f", Types.DoubleType.get()), | ||
| required(3, "dec", Types.DecimalType.of(11, 2))); | ||
|
|
||
| Record projected = | ||
| writeAndRead("type_promotion", writeSchema, promotedSchema, Lists.newArrayList(record)) | ||
| .get(0); | ||
|
|
||
| Assert.assertEquals("int should be promoted to long", 42L, projected.getField("id")); | ||
| Assert.assertEquals( | ||
| "float should be promoted to double", 1.5d, (double) projected.getField("f"), 0.0d); | ||
| Assert.assertEquals( | ||
| "decimal precision should widen", new BigDecimal("123.45"), projected.getField("dec")); | ||
| } | ||
|
|
||
| @Test | ||
| public void testReorderedProjectionValues() throws IOException { | ||
| Schema writeSchema = | ||
| new Schema( | ||
| required(1, "a", Types.LongType.get()), | ||
| optional(2, "b", Types.StringType.get()), | ||
| required(3, "c", Types.IntegerType.get())); | ||
|
|
||
| Record record = GenericRecord.create(writeSchema.asStruct()); | ||
| record.setField("a", 10L); | ||
| record.setField("b", "hello"); | ||
| record.setField("c", 7); | ||
|
|
||
| // Read schema requests the fields in a different order than the file's physical column order. | ||
| Schema reordered = | ||
| new Schema( | ||
| required(3, "c", Types.IntegerType.get()), | ||
| required(1, "a", Types.LongType.get()), | ||
| optional(2, "b", Types.StringType.get())); | ||
|
|
||
| Record projected = | ||
| writeAndRead("reordered", writeSchema, reordered, Lists.newArrayList(record)).get(0); | ||
|
|
||
| Assert.assertEquals("c should bind by id", 7, projected.getField("c")); | ||
| Assert.assertEquals("a should bind by id", 10L, projected.getField("a")); | ||
| Assert.assertEquals("b should bind by id", "hello", projected.getField("b").toString()); | ||
| } | ||
|
|
||
| @Test | ||
| public void testMetadataColumns() throws IOException { | ||
| Schema writeSchema = | ||
| new Schema( | ||
| required(1, "id", Types.LongType.get()), optional(2, "data", Types.StringType.get())); | ||
|
|
||
| List<Record> records = Lists.newArrayList(); | ||
| for (long i = 0; i < 5; i++) { | ||
| Record record = GenericRecord.create(writeSchema.asStruct()); | ||
| record.setField("id", i); | ||
| record.setField("data", "row" + i); | ||
| records.add(record); | ||
| } | ||
|
|
||
| Schema readSchema = | ||
| new Schema( | ||
| required(1, "id", Types.LongType.get()), | ||
| optional(2, "data", Types.StringType.get()), | ||
| MetadataColumns.ROW_POSITION, | ||
| MetadataColumns.IS_DELETED); | ||
|
|
||
| List<Record> projected = writeAndRead("metadata_columns", writeSchema, readSchema, records); | ||
|
|
||
| Assert.assertEquals(5, projected.size()); | ||
| for (int i = 0; i < projected.size(); i++) { | ||
| Record record = projected.get(i); | ||
| Assert.assertEquals("id should read from file", (long) i, record.getField("id")); | ||
| Assert.assertEquals( | ||
| "data should read from file", "row" + i, record.getField("data").toString()); | ||
| Assert.assertEquals( | ||
| "_pos should be the row position", | ||
| (long) i, | ||
| record.getField(MetadataColumns.ROW_POSITION.name())); | ||
| Assert.assertEquals( | ||
| "_deleted should be false", false, record.getField(MetadataColumns.IS_DELETED.name())); | ||
| } | ||
| } | ||
|
|
||
| @Test | ||
| public void testIdToConstant() throws IOException { | ||
| Schema writeSchema = | ||
| new Schema( | ||
| required(1, "id", Types.LongType.get()), optional(2, "data", Types.StringType.get())); | ||
|
|
||
| Record record = GenericRecord.create(writeSchema.asStruct()); | ||
| record.setField("id", 1L); | ||
| record.setField("data", "value"); | ||
|
|
||
| // Read schema adds an identity-partition-style constant field (id 3) not present in the file. | ||
| Schema readSchema = | ||
| new Schema( | ||
| required(1, "id", Types.LongType.get()), | ||
| optional(2, "data", Types.StringType.get()), | ||
| optional(3, "part", Types.StringType.get())); | ||
|
|
||
| Record projected = | ||
| writeAndRead( | ||
| "id_to_constant", | ||
| writeSchema, | ||
| readSchema, | ||
| Lists.newArrayList(record), | ||
| ImmutableMap.of(3, "constant-partition")) | ||
| .get(0); | ||
|
|
||
| Assert.assertEquals("id should read from file", 1L, projected.getField("id")); | ||
| Assert.assertEquals( | ||
| "data should read from file", "value", projected.getField("data").toString()); | ||
| Assert.assertEquals( | ||
| "part should resolve from idToConstant", "constant-partition", projected.getField("part")); | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -76,7 +76,7 @@ public OrcValueReader<?> record( | |
| TypeDescription record, | ||
| List<String> names, | ||
| List<OrcValueReader<?>> fields) { | ||
| return GenericOrcReaders.struct(fields, expected, idToConstant); | ||
| return GenericOrcReaders.struct(record, fields, expected, idToConstant); | ||
|
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Backported exactly. Visitor call site now passes the |
||
| } | ||
|
|
||
| @Override | ||
|
|
@@ -96,6 +96,10 @@ public OrcValueReader<?> map( | |
|
|
||
| @Override | ||
| public OrcValueReader<?> primitive(Type.PrimitiveType iPrimitive, TypeDescription primitive) { | ||
| if (iPrimitive == null) { | ||
| return null; | ||
| } | ||
|
|
||
| switch (primitive.getCategory()) { | ||
| case BOOLEAN: | ||
| return OrcValueReaders.booleans(); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -39,6 +39,7 @@ | |
| import org.apache.iceberg.types.Types; | ||
| import org.apache.iceberg.util.DateTimeUtil; | ||
| import org.apache.iceberg.util.UUIDUtil; | ||
| import org.apache.orc.TypeDescription; | ||
| import org.apache.orc.storage.ql.exec.vector.BytesColumnVector; | ||
| import org.apache.orc.storage.ql.exec.vector.ColumnVector; | ||
| import org.apache.orc.storage.ql.exec.vector.DecimalColumnVector; | ||
|
|
@@ -51,11 +52,25 @@ public class GenericOrcReaders { | |
|
|
||
| private GenericOrcReaders() {} | ||
|
|
||
| /** | ||
| * @deprecated Use {@link #struct(TypeDescription, List, Types.StructType, Map)} instead. This | ||
| * method uses position-based binding which may cause field misalignment in MOR and lineage | ||
| * scenarios. | ||
| */ | ||
| @Deprecated | ||
|
Comment on lines
+55
to
+60
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Backported with adaptations. The 3-arg overload is retained and marked |
||
| public static OrcValueReader<Record> struct( | ||
| List<OrcValueReader<?>> readers, Types.StructType struct, Map<Integer, ?> idToConstant) { | ||
| return new StructReader(readers, struct, idToConstant); | ||
| } | ||
|
|
||
| public static OrcValueReader<Record> struct( | ||
| TypeDescription orcType, | ||
| List<OrcValueReader<?>> readers, | ||
| Types.StructType struct, | ||
| Map<Integer, ?> idToConstant) { | ||
| return new StructReader(orcType, readers, struct, idToConstant); | ||
| } | ||
|
|
||
| public static OrcValueReader<List<?>> array(OrcValueReader<?> elementReader) { | ||
| return new ListReader(elementReader); | ||
| } | ||
|
|
@@ -204,12 +219,27 @@ public ByteBuffer nonNullRead(ColumnVector vector, int row) { | |
| private static class StructReader extends OrcValueReaders.StructReader<Record> { | ||
| private final GenericRecord template; | ||
|
|
||
| /** | ||
| * @deprecated Use {@link #StructReader(TypeDescription, List, Types.StructType, Map)} instead. | ||
| * This constructor uses position-based binding which may cause field misalignment in MOR | ||
| * and lineage scenarios. | ||
| */ | ||
| @Deprecated | ||
| protected StructReader( | ||
| List<OrcValueReader<?>> readers, | ||
| Types.StructType structType, | ||
| Map<Integer, ?> idToConstant) { | ||
| super(readers, structType, idToConstant); | ||
| this.template = structType != null ? GenericRecord.create(structType) : null; | ||
| this.template = GenericRecord.create(structType); | ||
| } | ||
|
Comment on lines
+228
to
+234
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Backported with adaptations. 4-arg subclass constructor from upstream. Adaptation: line 227 keeps this fork's |
||
|
|
||
| protected StructReader( | ||
| TypeDescription orcType, | ||
| List<OrcValueReader<?>> readers, | ||
| Types.StructType structType, | ||
| Map<Integer, ?> idToConstant) { | ||
| super(orcType, readers, structType, idToConstant); | ||
| this.template = GenericRecord.create(structType); | ||
| } | ||
|
|
||
| @Override | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Net new, using inspiration from shipped tests. Not part of the upstream commit. The metadata-column approach (project without metadata fields, hand the reader the full schema incl.
_pos/_deleted) is modeled on the shippedspark/v3.1/.../TestSparkOrcReadMetadataColumns.java; the write/project/read harness followsTestGenericReadProjection/TestGenericData. Covers type promotion, reordered projection, metadata columns, andidToConstant.