sunchao commented on code in PR #6753:
URL: https://github.com/apache/datafusion-comet/pull/6753#discussion_r4210417020
##########
spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala:
##########
@@ -45,4 +51,94 @@ object CometScanUtils {
case _ => false
}
}
+
+ /**
+ * Bin packs `files` with `pack` separately for each object store `storeKey`
names, so no
+ * partition holds files from two stores, and numbers the partitions from 0.
Files that share
+ * one store are packed exactly as `pack` packs them.
+ */
+ def packFilesPerStore[K](files: Seq[PartitionedFile], storeKey:
PartitionedFile => K)(
+ pack: Seq[PartitionedFile] => Seq[FilePartition]): Seq[FilePartition] = {
+ val keyed = files.map(file => (storeKey(file), file))
+ val stores = keyed.map(_._1).distinct
+ if (stores.size <= 1) {
+ pack(files)
+ } else {
+ stores
+ .flatMap(store => pack(keyed.collect { case (key, file) if key ==
store => file }))
Review Comment:
[P2] [P2] Group files once before packing each store. `keyed.collect` scans
the complete file list for every distinct store, making this driver-side step
O(files × stores). This also slows previously correct scans whose files already
occupy separate partitions. With 100,000 files across 1,000 stores, the
exact-head helper took about 2.97 seconds in a bounded probe, versus 97 ms for
a single grouping pass producing identical partitions. Please collect files
into groups once, preserving first-seen store order and file order, then invoke
`pack` for each group.
Evidence: Compiled the exact-head `NativeConfig` and changed
`CometScanUtils` methods against Spark 4.1.3. Constructed 100,000
`PartitionedFile`s with paths `s3a://bucket-${i % stores}/f$i` and
length/file_size 10,000. Invoked Spark's actual packing implementation with
maxSplitBytes and openCostInBytes both 134,217,728, yielding one file per
partition before and after the PR. After two warmups, median-of-three timings
were: 100 stores, base 12.76 ms / head 528.17 ms / grouping once 94.02 ms;
1,000 stores, base 4.91 ms / head 2973.80 ms / grouping once 96.60 ms. Asserted
identical indexed file layouts between the head and single-grouping
implementation. Probe and results:
`/tmp/comet6753-exact-head-7bx72sfj/probe/ScalingProbe.scala` and
`/tmp/comet6753-exact-head-7bx72sfj/ScalingProbe-isolated.log`.
--
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]