Skip to content

Blocking JVM calls from native plans hold Tokio workers, delaying other plans and I/O #6293

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions