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.
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:CometUdfBridge.evaluate(native/spark-expr/src/jvm_udf/mod.rs:221), once per batchCometS3CredentialProvider.getCredentialsForPath, once per S3 request (credential_bridge.rs:361), andgetPolicyLocationswhen a location-scoped store refreshesCometFileKeyUnwrapper.getKeyfor encrypted Parquet, which can go to the KMS on a cold cacheCelebornShufflePartitionPusherWhile 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
acquireMemorywas 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_ontook 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
acquireMemoryintokio::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 uporg/apache/comet/parquet/CometFileKeyUnwrapperby 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 onextraClassPath, as the installation guide recommends. The method ID could be resolved once inJVMClasses, like the others.