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]

Reply via email to