Using Comet with Apache Celeborn#
Comet accelerates query processing. Apache Celeborn provides remote storage for shuffle data.
Support Status#
Comet’s native shuffle through Celeborn is unavailable
with the currently released Celeborn 0.6.x and 0.7.x clients.
Setting spark.comet.shuffle.mode=native does not change that.
You can still use Comet to accelerate supported scans, filters, and other query operators. Shuffle is handled by Celeborn’s existing Spark integration.
Setup#
Start with a running Celeborn service and a compatible Comet installation. Celeborn is an optional application dependency that Comet does not bundle. Supply both the Comet JAR and a Celeborn Spark client matching the application’s Spark and Scala versions on the startup classpaths of the driver and all executors.
For Spark 3, the shaded client coordinates in Maven Central are
org.apache.celeborn:celeborn-client-spark-3-shaded_<scala-binary-version>:<celeborn-version>.
For example, org.apache.celeborn:celeborn-client-spark-3-shaded_2.12:0.7.0 supplies
celeborn-client-spark-3-shaded_2.12-0.7.0.jar for Spark 3 / Scala 2.12.
The example below uses Spark 3.5, Scala 2.12, and Celeborn 0.7.0:
Download the matching Comet JAR using the installation guide and the shaded Celeborn client from Maven Central.
Install both JARs in
$SPARK_HOME/jarson the driver and every executor, or include them in the Spark image used by those processes. For JARs installed elsewhere, setspark.driver.extraClassPathandspark.executor.extraClassPathto their paths before startup. Those paths must exist on the corresponding machines.Start Spark with the settings below, replacing the Celeborn master endpoints with your deployment’s values.
Spark 3.5 loads the shuffle manager before the executor’s user-JAR classloader is initialized.
Supplying only --jars or --packages does not ensure that the manager and Celeborn client
are available at that point; use the startup classpaths above.
This example enables Comet with Spark shuffle backed by Celeborn, as described in Support Status.
$SPARK_HOME/bin/spark-shell \
--conf spark.plugins=org.apache.spark.CometPlugin \
--conf spark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometCelebornShuffleManager \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
--conf spark.celeborn.master.endpoints=celeborn-master-1:9097,celeborn-master-2:9097 \
--conf spark.comet.exec.enabled=true \
--conf spark.comet.shuffle.enabled=true \
--conf spark.comet.shuffle.mode=auto \
--conf spark.comet.explain.fallback.enabled=true \
--conf spark.memory.offHeap.enabled=true \
--conf spark.memory.offHeap.size=2g
Celeborn’s shaded client requires the Kryo serializer; see the Celeborn deployment guide. For another Spark or Scala version, select matching Comet and Celeborn artifacts rather than reusing the Spark 3 / Scala 2.12 example.
Set the shuffle manager and Celeborn service configuration before creating the Spark context. Keep your deployment’s existing Celeborn authentication, storage, and recovery settings.
Verifying the Shuffle Path#
Run a query containing a shuffle, then inspect its executed plan in the Spark SQL UI or with
df.explain("formatted"). With AQE, inspect the final plan after an action completes. Supported operators can appear as Comet nodes, while the shuffle appears as a plainExchange. This is the expected result with the currently released Celeborn 0.6.x and 0.7.x clients.Inspect shuffle read/write bytes, records, and time in the Spark UI. Spark’s remote-read byte counters alone cannot confirm that Celeborn stored the data: they also count local shuffle files fetched from another executor. Celeborn’s fallback policy may select local Spark shuffle for some exchanges; use the check below to identify those shuffles.
If an operator you expected Comet to accelerate remains on Spark, use
spark.comet.explain.fallback.enabled=true to see the reasons in the driver log. See
Understanding Comet Plans for details.
Checking for Local Fallback#
For the Spark 3.5 / Celeborn 0.7.0 setup above, enable INFO logging for
org.apache.spark.scheduler.DAGScheduler on the driver before running the query, then:
Open the query in the Spark UI’s SQL tab, follow its associated jobs, and identify the shuffle-writing stage you want to check.
Find that stage in the driver log. For example,
Submitting ShuffleMapStage 5 (MapPartitionsRDD[17] ...)identifies its input RDD as 17. Find the correspondingRegistering RDD 17 (...) as input to shuffle 3message. In this example, stage 5 writes shuffle 3. Match the IDs within the same application; the messages need not be adjacent.Check Celeborn’s driver logs for
Fallback to vanilla Spark SortShuffleManager for shuffle: 3. This confirms that shuffle 3 selected local Spark shuffle. If fallback occurs with dynamic allocation enabled and no external shuffle service, Celeborn instead logs an ERROR containingfallback to vanilla Spark SortShuffleManager for shuffle: 3.
The IDs above are examples. Use the shuffle ID you found in step 2 when checking Celeborn’s fallback messages. The absence of a fallback message does not prove that data was stored in Celeborn.
Troubleshooting#
Symptom |
What to check |
|---|---|
Celeborn classes cannot be loaded when the application starts |
Supply the matching shaded client on both driver and executor startup classpaths. Comet does not bundle it. |
The plan contains |
This is expected with the currently released Celeborn 0.6.x and 0.7.x clients. Shuffle uses Celeborn’s existing Spark integration; other supported operators can still run in Comet. |
Setting |
Native shuffle is disabled with these clients. Changing this setting does not enable it. |
No Comet operators appear in the plan |
Check that the Comet plugin and native library loaded, that Comet execution is enabled, and that the query uses supported operators. Inspect the driver fallback explanations. |
Shuffle data is stored locally instead of in Celeborn |
Check Celeborn’s fallback policy, partition-count threshold, worker availability, and quota. An effective |
Changing the shuffle manager in the SQL session has no effect |
Set the manager before creating the Spark context. Restart the application to change it. |