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

   ### Describe the bug
   
   A native plan with no JVM input is polled on a Tokio worker 
(`jni_api.rs:1152-1189`), and several things it does call into the JVM 
synchronously from inside that poll:
   
   - JVM scalar UDFs, through `CometUdfBridge.evaluate` 
(`native/spark-expr/src/jvm_udf/mod.rs:221`), once per batch
   - `CometS3CredentialProvider.getCredentialsForPath`, once per S3 request 
(`credential_bridge.rs:361`), and `getPolicyLocations` when a location-scoped 
store refreshes
   - `CometFileKeyUnwrapper.getKey` for encrypted Parquet, which can go to the 
KMS on a cold cache
   - the executor-wide push admission wait in `CelebornShufflePartitionPusher`
   - the libhdfs NameNode calls that opendal's HDFS service makes inline in its 
async functions
   
   While one of these blocks, its worker runs nothing else. With one worker per 
task slot, that mostly costs the plan the overlap of its own I/O and compute. 
It has two wider effects.
   
   **Other plans wait.** When the runtime has fewer workers than plans ready to 
run, which is what the standalone fallback in #6292 produces, a blocked worker 
holds up other tasks' plans. Spark's `acquireMemory` was one of these calls 
until #6261: with a single worker it deadlocked, because the task holding the 
memory needed the worker the waiting task was blocking.
   
   **I/O stalls on the task threads.** Only Tokio workers drive the I/O and 
timer driver. A Spark task thread in `Handle::block_on`, which is how every 
plan with a JVM input runs, gets no timer or socket wake-ups while all workers 
are busy. That fits the worker-starvation explanation suggested in #6124 for 
the 10 s OpenDAL timeout.
   
   ### Steps to reproduce
   
   The I/O stall reproduces with tokio 1.53 alone. With one worker stuck in a 2 
s poll, a 10 ms sleep in another thread's `Handle::block_on` took 1.9 s, and so 
did a socket read whose peer wrote after 300 ms. With two workers both busy the 
result was the same.
   
   The Comet-level effect of a blocked worker is the deadlock in #6292, where 
the blocking call was `acquireMemory`. The calls above hold a worker the same 
way, but I haven't reproduced each of them in Comet.
   
   ### Expected behavior
   
   A JVM call that can block does not hold a Tokio worker while it blocks.
   
   ### Additional context
   
   #6261 wraps Spark's `acquireMemory` in `tokio::task::block_in_place`, which 
on a worker hands the worker's other tasks to another thread while the call 
blocks. On a Spark task thread it only steps out of the runtime context. #6261 
measured 0.1 to 2.4 µs per call, so the same wrapper around the calls above 
looks cheap.
   
   Related, though it is about class loading rather than blocking: 
`CometKeyRetriever::new` (`encryption_support.rs:91-110`) looks up 
`org/apache/comet/parquet/CometFileKeyUnwrapper` by name for every file, on 
whatever thread polls the scan. A Tokio worker has a null context class loader, 
so the lookup falls back to the system class loader. That works when Comet is 
on `extraClassPath`, as the installation guide recommends. The method ID could 
be resolved once in `JVMClasses`, like the others.
   


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