gaborkaszab commented on code in PR #18393: URL: https://github.com/apache/iceberg/pull/18393#discussion_r4193297289
########## core/src/main/java/org/apache/iceberg/formats/VectorizedStitchingIterable.java: ########## @@ -0,0 +1,198 @@ +/* + * 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.formats; + +import java.util.Collections; +import java.util.List; +import java.util.NoSuchElementException; +import org.apache.iceberg.io.CloseableGroup; +import org.apache.iceberg.io.CloseableIterable; +import org.apache.iceberg.io.CloseableIterator; +import org.apache.iceberg.io.SkippingCloseableIterator; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; + +/** + * Combines the batches read from a data file and its column files into batches of the requested + * projection. + * + * <p>Rows are matched by their position in the file, so only rows read by every part are produced. + * The parts may be read in batches with different boundaries, so each batch produced holds the rows + * that the current batches of all parts have in common. + * + * @param <D> the type of the batches that are combined + */ +class VectorizedStitchingIterable<D> extends CloseableGroup implements CloseableIterable<D> { + private final List<CloseableIterable<D>> parts; + private final VectorizedStitcher<D> stitcher; + + VectorizedStitchingIterable(List<CloseableIterable<D>> parts, VectorizedStitcher<D> stitcher) { + Preconditions.checkArgument(parts != null, "Invalid parts: null"); + Preconditions.checkArgument(parts.size() > 1, "Invalid parts: %s (must be > 1)", parts.size()); + Preconditions.checkArgument(stitcher != null, "Invalid stitcher: null"); + + this.parts = ImmutableList.copyOf(parts); + this.stitcher = stitcher; + + this.parts.forEach(this::addCloseable); + } + + @Override + public CloseableIterator<D> iterator() { + List<PartBatch<D>> partBatches = Lists.newArrayListWithCapacity(parts.size()); + for (int split = 0; split < parts.size(); split += 1) { + CloseableIterator<D> iterator = parts.get(split).iterator(); + if (!(iterator instanceof SkippingCloseableIterator<D> skipping)) { + throw new UnsupportedOperationException( + String.format("Cannot stitch vertical split %s: row positions are not tracked", split)); + } + + partBatches.add(new PartBatch<>(split, skipping, stitcher)); + } + + StitchingIterator<D> iterator = new StitchingIterator<>(partBatches, stitcher); + addCloseable(iterator); + return iterator; + } + + /** The current batch of a part and the range of row positions it holds. */ + private static class PartBatch<D> { + private final int split; Review Comment: `split` only seems to the used in an error message. I'm not sure how much information t gives to see that we were unable to read the 4th vertical split of the read. I'd rather remove it if there is no other intent. ########## parquet/src/main/java/org/apache/iceberg/parquet/VectorizedParquetReader.java: ########## @@ -139,23 +145,97 @@ public T next() { advance(); } - // batchSize is an integer, so casting to integer is safe - int numValuesToRead = (int) Math.min(nextRowGroupStart - valuesRead, batchSize); + return read(); + } + + /** Returns the position in the file of the first row of the next batch. */ + @Override + public long position() { + if (!hasNext()) { + throw new NoSuchElementException("No more rows"); + } + + if (valuesRead >= nextRowGroupStart) { + return rowIndexOffset(rowGroups.get(nextReadableRowGroup())); + } + + rowIndexOffset(rowGroups.get(nextRowGroup - 1)); + return nextPosition; + } + + /** + * Skips the batches that end at or before a position. Review Comment: If a position is in a row group that we filter, shouldn't we advance to the first valid option after the row group? That is what the row based implementation does ########## parquet/src/main/java/org/apache/iceberg/parquet/VectorizedParquetReader.java: ########## @@ -139,23 +145,97 @@ public T next() { advance(); } - // batchSize is an integer, so casting to integer is safe - int numValuesToRead = (int) Math.min(nextRowGroupStart - valuesRead, batchSize); + return read(); + } + + /** Returns the position in the file of the first row of the next batch. */ + @Override + public long position() { + if (!hasNext()) { + throw new NoSuchElementException("No more rows"); + } + + if (valuesRead >= nextRowGroupStart) { + return rowIndexOffset(rowGroups.get(nextReadableRowGroup())); + } + + rowIndexOffset(rowGroups.get(nextRowGroup - 1)); + return nextPosition; + } + + /** + * Skips the batches that end at or before a position. + * + * <p>Vectorized readers read whole batches, so the batch that holds the position is not split Review Comment: Hmm, I'm not entirely familiar with the internals of the batched readers, but this seems suspicious for me. Here we either stop before the position if the position in the next batch, or we stop after the position if the row group where the position is is filtered, or we stop at the position only if it's right at the beginning of a batch. If it's in the batch, can't we "move the needle" to the position and get the next batch starting from there? ########## core/src/main/java/org/apache/iceberg/formats/DataFileReadBuilder.java: ########## @@ -183,7 +183,12 @@ public CloseableIterable<D> build() { partReaders.add(builder.build()); }); - return new RowAlignedStitchingIterable<>(partReaders, stitcherBuilder.build(projection, parts)); + Stitcher<D> stitcher = stitcherBuilder.build(projection, parts); + if (stitcher instanceof VectorizedStitcher<D> vectorized) { Review Comment: I'm not sure we need a separate stitching utterable for vectorized reads. The intent was to have a general enough implementation in RowAlignedStitchingIterable where a row-based read is basically a batch-size=1 generalization but the algorithm should cover both. Currently, `RowAlignedStitchingIterable` (note, the name reflects that the files are row aligned and not that it is for row-based reads) can only read row-based, but I think the algorithm could be generalized to handle both cases. ########## core/src/main/java/org/apache/iceberg/formats/VectorizedStitchingIterable.java: ########## @@ -0,0 +1,198 @@ +/* + * 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.formats; + +import java.util.Collections; +import java.util.List; +import java.util.NoSuchElementException; +import org.apache.iceberg.io.CloseableGroup; +import org.apache.iceberg.io.CloseableIterable; +import org.apache.iceberg.io.CloseableIterator; +import org.apache.iceberg.io.SkippingCloseableIterator; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; + +/** + * Combines the batches read from a data file and its column files into batches of the requested + * projection. + * + * <p>Rows are matched by their position in the file, so only rows read by every part are produced. + * The parts may be read in batches with different boundaries, so each batch produced holds the rows + * that the current batches of all parts have in common. + * + * @param <D> the type of the batches that are combined + */ +class VectorizedStitchingIterable<D> extends CloseableGroup implements CloseableIterable<D> { + private final List<CloseableIterable<D>> parts; + private final VectorizedStitcher<D> stitcher; + + VectorizedStitchingIterable(List<CloseableIterable<D>> parts, VectorizedStitcher<D> stitcher) { + Preconditions.checkArgument(parts != null, "Invalid parts: null"); + Preconditions.checkArgument(parts.size() > 1, "Invalid parts: %s (must be > 1)", parts.size()); + Preconditions.checkArgument(stitcher != null, "Invalid stitcher: null"); + + this.parts = ImmutableList.copyOf(parts); + this.stitcher = stitcher; + + this.parts.forEach(this::addCloseable); + } + + @Override + public CloseableIterator<D> iterator() { + List<PartBatch<D>> partBatches = Lists.newArrayListWithCapacity(parts.size()); + for (int split = 0; split < parts.size(); split += 1) { + CloseableIterator<D> iterator = parts.get(split).iterator(); + if (!(iterator instanceof SkippingCloseableIterator<D> skipping)) { + throw new UnsupportedOperationException( + String.format("Cannot stitch vertical split %s: row positions are not tracked", split)); + } + + partBatches.add(new PartBatch<>(split, skipping, stitcher)); + } + + StitchingIterator<D> iterator = new StitchingIterator<>(partBatches, stitcher); + addCloseable(iterator); + return iterator; + } + + /** The current batch of a part and the range of row positions it holds. */ + private static class PartBatch<D> { + private final int split; + private final SkippingCloseableIterator<D> iterator; + private final VectorizedStitcher<D> stitcher; Review Comment: I don't think logically a Stitcher should be part of a `PartBatch`. It should sit on top of the batches and coordinate how to stitch them. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
