anuragmantri commented on code in PR #18324: URL: https://github.com/apache/iceberg/pull/18324#discussion_r4192147083
########## core/src/main/java/org/apache/iceberg/VerticalSplitProjectionPlanner.java: ########## @@ -0,0 +1,85 @@ +/* + * 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; + +import java.util.List; +import java.util.Map; +import java.util.Set; +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.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.TypeUtil; + +/** Plans the schema that each file of a vertically split data file is read with. */ +public class VerticalSplitProjectionPlanner { + private VerticalSplitProjectionPlanner() {} + + /** + * Prepares the schemas that can be used for each particular file as a projection for the reads. + * These schemas are prepared in a way to contain all the required fields to perform a full read + * including stitching vertical splits. + * + * @param file the data file to read, which may reference column files + * @param projection the schema the read produces + * @return a schema per file location to read it with, the base data file first + */ + public static Map<String, Schema> plan(DataFile file, Schema projection) { + Preconditions.checkArgument(file != null, "Invalid data file: null"); + Preconditions.checkArgument(projection != null, "Invalid projection: null"); + + List<ColumnFile> columnFiles = file.columnFiles(); + if (columnFiles == null || columnFiles.isEmpty()) { + return ImmutableMap.of(file.location(), projection); + } + + Set<Integer> projectedIdsToAssign = Sets.newLinkedHashSet(TypeUtil.getProjectedIds(projection)); Review Comment: Column files with nested types are planned incorrectly. The design doc says a column file that updates a nested field list only the top-level field ID. The planner matches field_ids against `TypeUtil.getProjectedIds(projection)`, which includes struct IDs and leaf IDs but not list or map IDs. I reproduced two problems through `GenericReader` implementation in #18388 1. A struct point `struct<x, y>` (ID 3, children 4 and 5), with a column file that has field_ids = [3]. The planner gives 3 to the column file, but 4 and 5 stay with the base file, and selecting them pulls point into the base projection as well. Both files then provide point, and the read fails with `Cannot stitch field 3: point ...: provided by more than one split`. 2. A List `tags list<string>` (ID 6, element 7), with a column file that has field_ids = `[6]`: ID 6 isn't among the projected IDs (only 7 is), so the column file is never read, and the read returns the base file's old tags values without an error.▎ Listing every nested ID `([3, 4, 5])` reads correctly today. Expanding each column file's IDs to all IDs under them before matching would fix both the above cases, for example `TypeUtil.getProjectedIds(TypeUtil.select(projection, columnFileIds))`, which gives `{3, 4, 5}` for the struct and `{7}` for the list. ########## core/src/main/java/org/apache/iceberg/VerticalSplitProjectionPlanner.java: ########## @@ -0,0 +1,85 @@ +/* + * 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; + +import java.util.List; +import java.util.Map; +import java.util.Set; +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.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.TypeUtil; + +/** Plans the schema that each file of a vertically split data file is read with. */ +public class VerticalSplitProjectionPlanner { + private VerticalSplitProjectionPlanner() {} + + /** + * Prepares the schemas that can be used for each particular file as a projection for the reads. + * These schemas are prepared in a way to contain all the required fields to perform a full read + * including stitching vertical splits. + * + * @param file the data file to read, which may reference column files + * @param projection the schema the read produces + * @return a schema per file location to read it with, the base data file first + */ + public static Map<String, Schema> plan(DataFile file, Schema projection) { + Preconditions.checkArgument(file != null, "Invalid data file: null"); + Preconditions.checkArgument(projection != null, "Invalid projection: null"); + + List<ColumnFile> columnFiles = file.columnFiles(); + if (columnFiles == null || columnFiles.isEmpty()) { + return ImmutableMap.of(file.location(), projection); + } + + Set<Integer> projectedIdsToAssign = Sets.newLinkedHashSet(TypeUtil.getProjectedIds(projection)); + Map<String, Schema> columnFilePlan = Maps.newLinkedHashMap(); + for (ColumnFile columnFile : columnFiles) { + Set<Integer> ids = Sets.newLinkedHashSet(); + for (Integer fieldId : columnFile.fieldIds()) { + if (projectedIdsToAssign.remove(fieldId)) { Review Comment: Should we have a `Preconditions.checkArgument` to check that the field id should not have been assigned to more than one column files? ########## core/src/main/java/org/apache/iceberg/formats/RowAlignedStitchingIterable.java: ########## @@ -0,0 +1,128 @@ +/* + * 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 row fragments read from a data file and its column files into the rows of the + * requested projection. + * + * <p>Fragments are matched by their position in the file, so only rows read by every part are + * produced. This keeps the parts aligned when a part reads only a range of the file or skips rows + * that cannot match a filter. + * + * @param <D> the type of the data records that are combined + */ +class RowAlignedStitchingIterable<D> extends CloseableGroup implements CloseableIterable<D> { + private final List<CloseableIterable<D>> parts; + private final Stitcher<D> stitcher; + + RowAlignedStitchingIterable(List<CloseableIterable<D>> parts, Stitcher<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<SkippingCloseableIterator<D>> iterators = 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)); + } + + iterators.add(skipping); + } + + return new StitchingIterator<>(iterators, stitcher); + } + + private static class StitchingIterator<D> extends CloseableGroup implements CloseableIterator<D> { + private final List<SkippingCloseableIterator<D>> iterators; + private final Stitcher<D> stitcher; + private final List<D> rows; + private boolean aligned = false; + + private StitchingIterator(List<SkippingCloseableIterator<D>> iterators, Stitcher<D> stitcher) { + this.iterators = iterators; + this.stitcher = stitcher; + this.rows = Lists.newArrayList(Collections.nCopies(iterators.size(), null)); + iterators.forEach(this::addCloseable); + } + + @Override + public boolean hasNext() { + // skip ahead until every split is at the same row + long target = 0; + int split = 0; + int alignedSplits = 0; + while (!aligned) { + SkippingCloseableIterator<D> iterator = iterators.get(split); + iterator.advanceTo(target); + if (!iterator.hasNext()) { Review Comment: Should we add a precondition here that the the column file reading should never stop when more base rows exist? (Dense representation invariant)? ########## core/src/main/java/org/apache/iceberg/formats/RowAlignedStitchingIterable.java: ########## @@ -0,0 +1,128 @@ +/* + * 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 row fragments read from a data file and its column files into the rows of the + * requested projection. + * + * <p>Fragments are matched by their position in the file, so only rows read by every part are + * produced. This keeps the parts aligned when a part reads only a range of the file or skips rows + * that cannot match a filter. + * + * @param <D> the type of the data records that are combined + */ +class RowAlignedStitchingIterable<D> extends CloseableGroup implements CloseableIterable<D> { + private final List<CloseableIterable<D>> parts; + private final Stitcher<D> stitcher; + + RowAlignedStitchingIterable(List<CloseableIterable<D>> parts, Stitcher<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<SkippingCloseableIterator<D>> iterators = 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)); + } + + iterators.add(skipping); + } + + return new StitchingIterator<>(iterators, stitcher); + } + + private static class StitchingIterator<D> extends CloseableGroup implements CloseableIterator<D> { + private final List<SkippingCloseableIterator<D>> iterators; + private final Stitcher<D> stitcher; + private final List<D> rows; + private boolean aligned = false; + + private StitchingIterator(List<SkippingCloseableIterator<D>> iterators, Stitcher<D> stitcher) { + this.iterators = iterators; + this.stitcher = stitcher; + this.rows = Lists.newArrayList(Collections.nCopies(iterators.size(), null)); + iterators.forEach(this::addCloseable); + } + + @Override + public boolean hasNext() { + // skip ahead until every split is at the same row + long target = 0; + int split = 0; + int alignedSplits = 0; + while (!aligned) { + SkippingCloseableIterator<D> iterator = iterators.get(split); + iterator.advanceTo(target); + if (!iterator.hasNext()) { + return false; + } + + long position = iterator.position(); + if (position == target) { + alignedSplits += 1; + } else { + target = position; + alignedSplits = 1; Review Comment: Another one I observed when prototyping #18393. This assumes `advanceTo(target)` never leaves an iterator before `target`. If it does, `target` moves backwards and two parts can keep resetting each other forever. Maybe add a `Preconditions.checkState(position >= target, ...)` here to fail fast? ########## core/src/main/java/org/apache/iceberg/formats/Stitcher.java: ########## @@ -0,0 +1,40 @@ +/* + * 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.List; + +/** + * Combines the vertical splits of a row into a single row. + * + * @param <D> the type of the data records that are combined + */ +public interface Stitcher<D> { + /** + * Combines one unit of each vertical split into a unit of the projection. + * + * <p>The list is reused by the caller, so implementations must not retain it. + */ + D stitch(List<D> parts, int count); + + /** Returns a range of rows of a unit; only supported by vectorized models. */ + default D slice(D unit, int offset, int count) { + throw new UnsupportedOperationException("Slicing is not supported for row based models"); + } Review Comment: `slice` has no caller in the row path, and its default throws. Would a separate interface for vectorized models be cleaner, with `slice` and a row count, so the read builder can choose a batch iterable by type? #18393 adds `VectorizedStitcher` that way, but I'm happy to fold the methods into `Stitcher` if you prefer. ########## api/src/main/java/org/apache/iceberg/io/SkippingCloseableIterator.java: ########## @@ -0,0 +1,44 @@ +/* + * 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.io; + +/** + * An iterator over the rows of a file that knows the position of each row in the file and can skip + * ahead to a position. + * + * @param <T> the type of the rows + */ +public interface SkippingCloseableIterator<T> extends CloseableIterator<T> { + /** + * Returns the position in the file of the row that the next call to {@link #next()} returns. + * + * @return the position of the next row + * @throws java.util.NoSuchElementException if there are no more rows + */ + long position(); + + /** + * Skips the rows before a position, so the next row returned is the first one at or after it. + * + * <p>Does nothing if the next row is already at or after the position. + * + * @param position the position in the file to skip to + */ + void advanceTo(long position); Review Comment: I prototyped batch stitching in #18393 . For batch readers, `advanceTo` can't stop in the middle of a batch: `VectorizedArrowReader.read` always reads a whole batch. Could the Javadoc define the unit? For example, `position()` is the position of the first row of the next unit (a row or a batch), and `advanceTo` may stop at a unit that starts before the target. ########## spark/v4.2/spark/src/test/java/org/apache/iceberg/ColumnFileTestHelpers.java: ########## @@ -0,0 +1,74 @@ +/* + * 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; + +import java.util.List; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; + +public final class ColumnFileTestHelpers { Review Comment: nit: could this live in core test sources instead of Spark's? #18388 moves it to core so the data module's tests can use it, and other modules already depend on core's test output. Moving it here would drop that rename from my PR. -- 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]
