Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
*/
package org.apache.iceberg.arrow.vectorized;

import java.util.List;
import org.apache.arrow.vector.FieldVector;
import org.apache.iceberg.MetadataColumns;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;
Expand Down Expand Up @@ -215,4 +216,35 @@ public VectorHolder valueHolder() {
return valueHolder;
}
}

public static class StructVectorHolder extends VectorHolder {
private final int numRows;
private final List<VectorHolder> childHolders;
private final NullabilityHolder structNulls;

public StructVectorHolder(
Types.NestedField icebergField,
int numRows,
List<VectorHolder> childHolders,
NullabilityHolder structNulls) {
super(icebergField);
this.numRows = numRows;
this.childHolders = childHolders;
this.structNulls = structNulls;
}

@Override
public int numValues() {
return numRows;
}

public List<VectorHolder> childHolders() {
return childHolders;
}

@Override
public NullabilityHolder nullabilityHolder() {
return structNulls;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
*/
package org.apache.iceberg.arrow.vectorized;

import java.util.List;
import java.util.Map;
import java.util.Optional;
import org.apache.arrow.memory.ArrowBuf;
Expand Down Expand Up @@ -47,6 +48,7 @@
import org.apache.iceberg.parquet.ParquetUtil;
import org.apache.iceberg.parquet.VectorizedReader;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.types.Types;
import org.apache.parquet.column.ColumnDescriptor;
import org.apache.parquet.column.Dictionary;
Expand Down Expand Up @@ -136,6 +138,14 @@ public void setBatchSize(int batchSize) {
this.vectorizedColumnIterator.setBatchSize(batchSize);
}

void registerStructPresence(NullabilityHolder structNulls, int structDefinitionLevel) {
vectorizedColumnIterator.addStructPresence(structNulls, structDefinitionLevel);
}

protected VectorizedArrowReader fileBackedLeaf() {
return vectorizedColumnIterator != null ? this : null;
}

@Override
public VectorHolder read(VectorHolder reuse, int numValsToRead) {
boolean dictEncoded = vectorizedColumnIterator.producesDictionaryEncodedVector();
Expand Down Expand Up @@ -386,6 +396,11 @@ public void setRowGroupInfo(PageReadStore source, Map<ColumnPath, ColumnChunkMet
!ParquetUtil.hasNonDictionaryPages(chunkMetaData));
}

boolean isColumnInRowGroup(Map<ColumnPath, ColumnChunkMetaData> metadata) {
return columnDescriptor != null
&& metadata.containsKey(ColumnPath.get(columnDescriptor.getPath()));
}

@Override
public void close() {
if (vec != null) {
Expand Down Expand Up @@ -1054,6 +1069,12 @@ public VectorizedVariantReader(
this.valueReader = valueReader;
}

@Override
protected VectorizedArrowReader fileBackedLeaf() {
VectorizedArrowReader metadataLeaf = metadataReader.fileBackedLeaf();
return metadataLeaf != null ? metadataLeaf : valueReader.fileBackedLeaf();
}

@Override
public VectorHolder read(VectorHolder reuse, int numValsToRead) {
VectorHolder reuseMetadata = null;
Expand Down Expand Up @@ -1093,4 +1114,120 @@ public String toString() {
return "VectorizedVariantReader";
}
}

static class VectorizedStructReader extends VectorizedArrowReader {
private final List<VectorizedReader<?>> childReaders;
private final int structDefinitionLevel;
private final VectorizedArrowReader presenceReader;
private NullabilityHolder structNulls;
private VectorHolder reusePresence;
private boolean presenceColumnInRowGroup;

VectorizedStructReader(
Types.NestedField icebergField,
List<VectorizedReader<?>> childReaders,
int structDefinitionLevel,
VectorizedArrowReader presenceReader) {
super(icebergField);
this.childReaders = childReaders;
this.structDefinitionLevel = structDefinitionLevel;
this.presenceReader = presenceReader;
}

@Override
protected VectorizedArrowReader fileBackedLeaf() {
for (VectorizedReader<?> child : childReaders) {
if (child instanceof VectorizedArrowReader) {
VectorizedArrowReader leaf = ((VectorizedArrowReader) child).fileBackedLeaf();
if (leaf != null) {
return leaf;
}
}
}

// expose the presence leaf so an ancestor reuses it instead of double-reading the column
return presenceReader;
}

@Override
public VectorHolder read(VectorHolder reuse, int numValsToRead) {
if (structNulls != null) {
structNulls.reset();
}

// when presenceColumnInRowGroup is false the struct is always present (structNulls left all
// not-null)
if (presenceReader != null && presenceColumnInRowGroup) {
// populate structNulls as a side effect of a value-column batch read; the values are unused
this.reusePresence = presenceReader.read(reusePresence, numValsToRead);
}

List<VectorHolder> reuseChildren = null;
if (reuse instanceof VectorHolder.StructVectorHolder) {
reuseChildren = ((VectorHolder.StructVectorHolder) reuse).childHolders();
}

List<VectorHolder> childHolders = Lists.newArrayListWithExpectedSize(childReaders.size());
for (int idx = 0; idx < childReaders.size(); idx++) {
VectorHolder reuseChild = reuseChildren == null ? null : reuseChildren.get(idx);
VectorizedArrowReader child = (VectorizedArrowReader) childReaders.get(idx);
childHolders.add(child.read(reuseChild, numValsToRead));
}

return new VectorHolder.StructVectorHolder(
icebergField(), numValsToRead, childHolders, structNulls);
}

@Override
public void setRowGroupInfo(
PageReadStore source, Map<ColumnPath, ColumnChunkMetaData> metadata) {
for (VectorizedReader<?> child : childReaders) {
child.setRowGroupInfo(source, metadata);
}

// a partition-constant presence column is absent from the row group, so the struct is present
this.presenceColumnInRowGroup =
presenceReader != null && presenceReader.isColumnInRowGroup(metadata);
if (presenceColumnInRowGroup) {
presenceReader.setRowGroupInfo(source, metadata);
}
}

@Override
public void setBatchSize(int batchSize) {
int resolvedBatchSize = (batchSize == 0) ? DEFAULT_BATCH_SIZE : batchSize;
for (VectorizedReader<?> child : childReaders) {
child.setBatchSize(resolvedBatchSize);
}

if (presenceReader != null) {
presenceReader.setBatchSize(resolvedBatchSize);
}

if (structDefinitionLevel > 0) {
// per-row presence from a file leaf under the struct (child or shared presence leaf)
VectorizedArrowReader presenceLeaf = fileBackedLeaf();
if (presenceLeaf != null) {
this.structNulls = new NullabilityHolder(resolvedBatchSize);
presenceLeaf.registerStructPresence(structNulls, structDefinitionLevel);
}
}
}

@Override
public void close() {
for (VectorizedReader<?> child : childReaders) {
child.close();
}

if (presenceReader != null) {
presenceReader.close();
}
}

@Override
public String toString() {
return "VectorizedStructReader(" + childReaders.size() + ")";
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.iceberg.Schema;
import org.apache.iceberg.arrow.ArrowAllocation;
import org.apache.iceberg.arrow.vectorized.VectorizedArrowReader.ConstantVectorReader;
import org.apache.iceberg.parquet.ParquetSchemaUtil;
import org.apache.iceberg.parquet.ParquetVariantVisitor;
import org.apache.iceberg.parquet.TypeWithSchemaVisitor;
import org.apache.iceberg.parquet.VectorizedReader;
Expand Down Expand Up @@ -102,27 +103,32 @@ protected VectorizedReaderBuilder(
@Override
public VectorizedReader<?> message(
Types.StructType expected, MessageType message, List<VectorizedReader<?>> fieldReaders) {
GroupType groupType = message.asGroupType();
Map<Integer, VectorizedReader<?>> readersById = Maps.newHashMap();
List<Type> fields = groupType.getFields();

IntStream.range(0, fields.size())
.filter(pos -> fields.get(pos).getId() != null)
.forEach(pos -> readersById.put(fields.get(pos).getId().intValue(), fieldReaders.get(pos)));

List<Types.NestedField> icebergFields =
expected != null ? expected.fields() : ImmutableList.of();
return vectorizedReader(
reorderFields(icebergFields, message.asGroupType().getFields(), fieldReaders));
}

List<VectorizedReader<?>> reorderedFields =
Lists.newArrayListWithExpectedSize(icebergFields.size());

for (Types.NestedField field : icebergFields) {
private List<VectorizedReader<?>> reorderFields(
List<Types.NestedField> expectedFields,
List<Type> parquetFields,
List<VectorizedReader<?>> fieldReaders) {
Map<Integer, VectorizedReader<?>> readersById = Maps.newHashMap();
IntStream.range(0, parquetFields.size())
.filter(pos -> parquetFields.get(pos).getId() != null)
.forEach(
pos ->
readersById.put(parquetFields.get(pos).getId().intValue(), fieldReaders.get(pos)));

List<VectorizedReader<?>> reordered = Lists.newArrayListWithExpectedSize(expectedFields.size());
for (Types.NestedField field : expectedFields) {
VectorizedReader<?> reader =
VectorizedArrowReader.replaceWithMetadataReader(
field, readersById.get(field.fieldId()), idToConstant, setArrowValidityVector);
reorderedFields.add(defaultReader(field, reader));
reordered.add(defaultReader(field, reader));
}
return vectorizedReader(reorderedFields);

return reordered;
}

private VectorizedReader<?> defaultReader(Types.NestedField field, VectorizedReader<?> reader) {
Expand All @@ -148,11 +154,55 @@ protected VectorizedReader<?> vectorizedReader(List<VectorizedReader<?>> reorder
@Override
public VectorizedReader<?> struct(
Types.StructType expected, GroupType groupType, List<VectorizedReader<?>> fieldReaders) {
if (expected != null) {
throw new UnsupportedOperationException(
"Vectorized reads are not supported yet for struct fields");
if (expected == null) {
return null;
}

// no field ID / no matching Iceberg field: fall back like primitive()
if (groupType.getId() == null) {
return null;
}
return null;

Types.NestedField structField = icebergSchema.findField(groupType.getId().intValue());
if (structField == null) {
return null;
}

List<VectorizedReader<?>> reorderedFields =
reorderFields(expected.fields(), groupType.getFields(), fieldReaders);

int structDefinitionLevel = parquetSchema.getMaxDefinitionLevel(currentPath());

VectorizedArrowReader presenceReader = null;
// must agree with ParquetSchemaUtil.PresenceColumnSelector, which retains this presence column
if (structDefinitionLevel > 0 && !hasFileBackedLeaf(reorderedFields)) {
ColumnDescriptor presence =
ParquetSchemaUtil.selectPresenceColumn(parquetSchema, currentPath());
// list/map are not vectorized, so the presence leaf is always rep level 0 (iterator
// precondition)
if (presence != null && presence.getMaxRepetitionLevel() == 0) {
presenceReader =
new VectorizedArrowReader(
presence,
ParquetSchemaUtil.presenceField(presence),
rootAllocator,
setArrowValidityVector);
}
}

return new VectorizedArrowReader.VectorizedStructReader(
structField, reorderedFields, structDefinitionLevel, presenceReader);
}

private boolean hasFileBackedLeaf(List<VectorizedReader<?>> readers) {
for (VectorizedReader<?> child : readers) {
if (child instanceof VectorizedArrowReader
&& ((VectorizedArrowReader) child).fileBackedLeaf() != null) {
return true;
}
}

return false;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,11 @@ public VectorizedColumnIterator(

public void setBatchSize(int batchSize) {
this.batchSize = batchSize;
vectorizedPageIterator.clearStructPresences();
}

public void addStructPresence(NullabilityHolder structNulls, int structDefinitionLevel) {
vectorizedPageIterator.addStructPresence(structNulls, structDefinitionLevel);
}

public Dictionary setRowGroupInfo(PageReader store, boolean allPagesDictEncoded) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,13 @@
package org.apache.iceberg.arrow.vectorized.parquet;

import java.io.IOException;
import java.util.List;
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.IntVector;
import org.apache.iceberg.arrow.vectorized.NullabilityHolder;
import org.apache.iceberg.parquet.BasePageIterator;
import org.apache.iceberg.parquet.ParquetUtil;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.parquet.CorruptDeltaByteArrays;
import org.apache.parquet.bytes.ByteBufferInputStream;
import org.apache.parquet.bytes.BytesUtils;
Expand Down Expand Up @@ -58,6 +60,19 @@ private enum DictionaryDecodeMode {

private DictionaryDecodeMode dictionaryDecodeMode;

private final List<VectorizedParquetDefinitionLevelReader.StructPresence> structPresences =
Lists.newArrayList();

void addStructPresence(NullabilityHolder structNulls, int structDefinitionLevel) {
structPresences.add(
new VectorizedParquetDefinitionLevelReader.StructPresence(
structNulls, structDefinitionLevel));
}

void clearStructPresences() {
structPresences.clear();
}

public void setAllPagesDictEncoded(boolean allDictEncoded) {
this.allPagesDictEncoded = allDictEncoded;
}
Expand Down Expand Up @@ -148,6 +163,7 @@ protected void initDefinitionLevelsReader(
this.vectorizedDefinitionLevelReader =
new VectorizedParquetDefinitionLevelReader(
bitWidth, desc.getMaxDefinitionLevel(), setArrowValidityVector);
this.vectorizedDefinitionLevelReader.setStructPresences(structPresences);
this.vectorizedDefinitionLevelReader.initFromPage(triplesCount, in);
}

Expand All @@ -159,6 +175,7 @@ protected void initDefinitionLevelsReader(DataPageV2 dataPageV2, ColumnDescripto
this.vectorizedDefinitionLevelReader =
new VectorizedParquetDefinitionLevelReader(
bitWidth, desc.getMaxDefinitionLevel(), false, setArrowValidityVector);
this.vectorizedDefinitionLevelReader.setStructPresences(structPresences);
this.vectorizedDefinitionLevelReader.initFromPage(
dataPageV2.getValueCount(), dataPageV2.getDefinitionLevels().toInputStream());
}
Expand Down
Loading
Loading