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]
