andygrove opened a new issue, #6313:
URL: https://github.com/apache/datafusion-comet/issues/6313

   ### Describe the bug
   
   When a native block has no JVM input, meaning every leaf is a native scan, 
`executePlan` publishes the whole metric tree to the JVM after every output 
batch (the `batch_receiver` branch in `native/core/src/execution/jni_api.rs`). 
Blocks with a JVM input publish only after `spark.comet.metrics.updateInterval` 
has passed. So pure-native blocks ignore the setting, including a negative 
value, which the config doc says means "metrics will be updated upon task 
completion".
   
   This came in with #3553 (0.14.0), which moved pure-native blocks onto a 
channel. Before that, both kinds of block checked the interval.
   
   Each publish walks the plan, runs `aggregate_by_name` over every node's 
metrics, encodes a protobuf, and calls into the JVM, which decodes it and sets 
each `SQLMetric`. A native Parquet scan registers about 45 metrics for each 
file it opens, so the cost grows with the number of files a task has read.
   
   I timed `update_metrics` in place on a release build (Apple M3 Max, 
`local[1]`, 16M rows in 2,048 output batches from one task):
   
   | Plan                                               | Cost per publish | 
Per task |
   | -------------------------------------------------- | ---------------- | 
-------- |
   | Scan only, 1 file                                  | 19 µs            | 40 
ms    |
   | Scan + filter + project                            | 30 µs            | 62 
ms    |
   | Scan only, 64 files in one task (2,905 raw metrics) | 38 µs            | 
78 ms    |
   | Scan + filter + project, `noop` write              | 26 µs            | 54 
ms    |
   
   Switching that branch to the interval check gave these medians of 7 runs:
   
   | Plan                                  | Wall time      | Process CPU time |
   | ------------------------------------- | -------------- | ---------------- |
   | Scan only, 1 file                     | 174 → 143 ms   | 343 → 209 ms     |
   | Scan + filter + project               | 714 → 716 ms   | 882 → 834 ms     |
   | Scan only, 64 files in one task       | 210 → 167 ms   | 313 → 220 ms     |
   | Scan + filter + project, `noop` write | 917 → 848 ms   | 1608 → 1505 ms   |
   
   The scan + filter + project wall time doesn't move because the producer is 
the bottleneck there, and the publish runs on the consumer thread. It still 
spends CPU that other tasks on the executor could use.
   
   This affects pure-native blocks whose output reaches the JVM batch by batch: 
a scan feeding a Spark write or a collect, a broadcast build side, JVM shuffle, 
or a fallback operator. A block that ends in a native shuffle write hands back 
only one batch, so it isn't affected.
   
   ### Steps to reproduce
   
   Set `spark.comet.metrics.updateInterval=-1` and 
`spark.comet.batchSize=1000`, and read a 10,000-row Parquet file in one task. 
Inside the task, read the native scan's `output_rows` metric after the first 
batch. It is already non-zero, when it should stay 0 until the iterator closes.
   
   ### Expected behavior
   
   Pure-native blocks publish on the configured interval, as blocks with a JVM 
input do, and `releasePlan` publishes the final values.
   
   ### Additional context
   
   Found while looking at #1381.
   


-- 
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