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

   ## Is your feature request related to a problem or challenge?
   
   On the executor, every task builds a brand-new `RuntimeEnv` from scratch. 
The default runtime producer (`ballista/executor/src/executor_process.rs`, 
`run_executor`) is:
   
   ```rust
   Arc::new(move |_| {
       let runtime_env = RuntimeEnvBuilder::new()
           .with_temp_file_path(wd.clone())
           .build()?;
       Ok(Arc::new(runtime_env))
   })
   ```
   
   This closure is invoked once per task (`get_task_definition` / 
`get_task_definition_vec` in `ballista/core/src/serde/scheduler/from_proto.rs`, 
and `execution_loop.rs`), so each `TaskContext` gets a fresh `DiskManager`, 
`CacheManager`, `ObjectStoreRegistry`, and memory pool. Nothing runtime-scoped 
is shared between concurrent tasks on the same executor — including tasks of 
the same job/stage.
   
   The optional `--memory-pool-size` wrapper 
(`wrap_runtime_producer_with_memory_pool`, #1624) does not change this: it 
still calls the inner producer (fresh env) per task and layers a **fresh 
per-task `FairSpillPool`** of `total_bytes / concurrent_tasks` on top.
   
   Consequences:
   
   1. **No cross-task memory arbitration (the main one).** Memory is statically 
partitioned per task slot. A task never sees the executor's real aggregate 
memory pressure and cannot borrow budget from idle sibling slots. Under 
partition skew — the case Ballista's scalability work most cares about — a hot 
task that could stay in-memory with more headroom is instead capped at one 
slot's slice and forced to spill, while neighboring slots sit idle. The default 
(no `--memory-pool-size`) is an unbounded pool per task, i.e. no accounting at 
all.
   2. **No object-store client reuse.** A fresh `ObjectStoreRegistry` per task 
means object-store clients (connection pools, credentials, TLS sessions) are 
rebuilt per task. Negligible for local Parquet, a real per-task tax against 
S3/GCS at scale.
   3. **No file metadata / statistics / listing cache reuse.** A fresh 
`CacheManager` per task means Parquet footers, file listings, and statistics 
are re-read by every task and every stage, even when many tasks on the same 
executor scan the same table.
   
   ## Describe the solution you'd like
   
   Allow runtime-scoped state to be shared executor-wide instead of rebuilt per 
task, at least for the pieces where sharing is safe, and evaluate sharing the 
memory pool:
   
   - **Shared `ObjectStoreRegistry` and `CacheManager`**: build these once at 
executor startup and reuse them across tasks (preserve them through the runtime 
producer rather than reconstructing). Low risk; these are designed to be shared.
   - **Executor-wide memory pool (evaluate)**: replace the per-task 
`FairSpillPool` slice with a single executor-wide pool shared by all concurrent 
tasks, so a hot task can expand into slack left by idle slots. This is a real 
trade-off — a shared pool gives better utilization under skew but weaker 
isolation (one greedy task can starve/OOM others), whereas static per-task 
slices give predictability. Worth measuring both on the SF10 harness with a 
skewed query, under AQE on and AQE off.
   
   An escape hatch already exists: 
`ExecutorProcessConfig::override_runtime_producer` lets an embedder inject a 
producer closing over shared singletons, so (1)/(2)/(3) are achievable today 
without core changes — but the default does not do it, and the memory-pool 
wrapper deliberately makes a fresh per-task pool.
   
   ## Describe alternatives you've considered
   
   - Keep the per-task isolation model and instead make the static memory slice 
smarter (e.g. weight by stage). Simpler, but still can't reclaim idle slots' 
budget for a hot task.
   - Share only the read-side caches (`ObjectStoreRegistry`, `CacheManager`) 
and leave the memory pool per-task. This captures the low-risk wins (2)/(3) 
without the isolation trade-off of a shared pool.
   
   ## Additional context
   
   Nothing runtime-scoped currently persists across tasks; the executor's 
`RuntimeEnv` is fully task-local. Related, but distinct: #1624 (per-task 
`--memory-pool-size`), #1578 (cache ballista clients on executor — gRPC 
clients, not object store), #1836 (coordinator OOM at SF100), #1952 (spill 
large shuffle fetches).
   


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