NoahKusaba opened a new issue, #2509:
URL: https://github.com/apache/datafusion-ballista/issues/2509

   ### Is your feature request related to a problem or challenge?
   
   The shuffle writers do their CPU work on threads that the executor's vcore 
accounting doesn't count.
   
   `utils::write_stream_to_disk` (`ballista/core/src/utils.rs`) forwards 
batches through a channel to a `spawn_blocking` task. That task holds a 
blocking-pool thread for the whole life of the output stream, and it does IPC 
encoding and LZ4/ZSTD compression as well as the file writes. The passthrough 
path in `ShuffleWriterExec` (`shuffle_writer.rs`) drains all K output 
partitions at once, so a task runs K of these threads. 
`range_shuffle::write_stream_to_ipc_file` follows the same pattern.
   
   The executor runs tasks on a `DedicatedExecutor` whose worker count equals 
`--vcores` (`execution_loop.rs`, `executor_server.rs`). Work passed to 
`spawn_blocking` runs on tokio's blocking pool instead, which isn't bounded by 
vcores (it defaults to up to 512 threads). So:
   
   - with K = 200, one task can hold 200 OS threads, all compressing data;
   - a task the scheduler assigned 4 vcores can use every core on the host, 
which breaks the vcore accounting the scheduler relies on and takes CPU from 
other tasks on the same executor.
   
   Measured with a small benchmark that calls `write_stream_to_disk` directly, 
on a 12-core machine with 200 partitions per task: a single 4-vcore task used 
about 10.7 cores and peaked at 205 OS threads. With 3 tasks on a 12-vcore 
runtime, the process peaked at 526 threads.
   
   There are also some synchronous file operations on async workers in the same 
path (`create_dir_all` before each write, the range-shuffle index write).
   
   ### Describe the solution you'd like
   
   Separate CPU work from I/O, and give every file operation in the shuffle 
path the same policy:
   
   - Encode and compress on the task's own runtime, so the work counts against 
its vcores.
   - Send the finished byte buffers to one bounded, executor-wide I/O pool that 
only does file operations.
   - Size that pool with an executor option (e.g. `--shuffle-io-threads`), 
install it on each task's `SessionConfig`, and fall back to a process-wide 
default when none is installed (tests, client-side execution).
   - Bound memory per output file: one write in flight, a limited queue of 
encoded bytes behind it, and backpressure on the producer once the queue is 
full.
   
   A prototype of this (hash/passthrough and range shuffle writers) keeps a 
4-vcore task at about 4.3 cores and 14 threads (down from 205). On a fully 
loaded 12-vcore runtime it peaks at 22 threads (down from 526), with the same 
or better wall time and about 10% less CPU with zstd.
   
   One case still needs work before a PR: when writes exceed the kernel's 
dirty-page limit and the disk becomes the bottleneck, the prototype is slower 
than the current code. The vcore threads fill the per-file queues and then idle 
behind throttled I/O threads, while the current code hides that wait behind 
hundreds of threads. Bounding the per-file queue by bytes instead of by batch 
count is the next thing to try.
   
   The sort-based shuffle writer has a related problem: its spill and final 
consolidated write run synchronously on async workers. That can move onto the 
same pool as a follow-up.
   
   ### Describe alternatives you've considered
   
   - Keep `spawn_blocking` but cap concurrency with a semaphore. That bounds 
the thread count, but compression still runs outside the vcore-sized runtime.
   - Async file I/O (`tokio::fs`). It still uses the blocking pool internally, 
and it doesn't address where the encode/compress CPU work runs.
   
   ### Additional context
   
   Follow-up to #1387, which moved the shuffle writer's synchronous I/O off the 
async workers via `spawn_blocking`. That fixed the blocked workers but left the 
CPU work, and one thread per output stream, on the blocking pool.
   


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