andygrove opened a new issue, #6512: URL: https://github.com/apache/datafusion-comet/issues/6512
### What is the problem the feature request solves? When Spark splits a Parquet file into several partitions, the native Parquet scan reads the same byte ranges as Spark: `get_partitioned_files` in `native/core/src/execution/planner.rs` hands each split's start and length straight to DataFusion. The two engines disagree about which split reads a row group that crosses a split boundary, though. DataFusion keeps a row group in the split that holds the first page of its first column (`row_group_in_range` in DataFusion's `row_group_filter.rs`; in 55.1 the check is inline in `prune_by_range`). parquet-mr keeps it in the split that holds its midpoint, which is that offset plus half the row group's compressed size (`ParquetMetadataConverter.filterFileMetaDataByMidpoint`). Every row is still read exactly once, but sometimes by a different split, and so in a different partition, than in Spark. That's the cause of two issues we know about. `_metadata.file_block_start` and `file_block_length` report the wrong split (#6505), which #6510 works around by falling back to Spark when a query reads either column. And Comet uses less of the parallelism Spark planned (#3817): when row groups are close to `spark.sql.files.maxPartitionBytes` in size, some splits that read a row group in Spark read nothing in Comet, and their neighbors read more. The scan tuning guide suggests lowering `maxPartitionBytes` to work around that. ### Describe the potential solution Make DataFusion's `row_group_in_range` use the midpoint rule. That's what DataFusion's docs already say it does: the rustdoc on `FileRange` says it scans the row groups "whose data 'midpoint' lies within the [start, end) byte offsets". apache/datafusion#1990 added `FileRange` with that behavior in 2022 by calling arrow-rs's `ReadOptionsBuilder::with_range`, which still uses the midpoint. apache/datafusion#2677 replaced that with a check on the first column chunk's `file_offset` when the scan moved to `ParquetRecordBatchStream`, and apache/datafusion#5997 changed it to the first page offset. The docs never changed. I haven't found a DataFusion issue for this yet. The change is a few lines, and every row group still lands in exactly one split. It should find the start of the row group the way parquet-mr does, as the smaller of the first column's dictionary and data page offsets. iceberg-rust instead adds up compressed sizes from offset 4, which drifts on files where parquet-java padded row groups out to an HDFS block boundary (up to 8 MB per row group by default). Once Comet picks up a DataFusion release with the change, we can revert the #6510 fallback, add a `_metadata.file_block_start` test that splits a file with several row groups, and rewrite the "Parquet Native Scans" section of the scan tuning guide. That section also says Spark assigns a row group by its start offset, which isn't right. Doing it in Comet alone would be harder. DataFusion's `ParquetAccessPlan` needs one entry per row group, so Comet would have to read every footer before building the plan, or wrap `ParquetSource`. ### Additional context The native Iceberg scan already uses the midpoint rule, since apache/iceberg-rust#2615 fixed #4590. cc @comphead, who looked into the parallelism side in #3817. -- 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]
