All Products
Search
Document Center

E-MapReduce:Celeborn

Last Updated:Aug 20, 2026

Celeborn processes intermediate data to enhance the stability, flexibility, and performance of big data engines. This topic explains how to use the Celeborn service.

Background information

Existing shuffle solutions have the following disadvantages:

  • In scenarios with large data volumes, shuffle writes can cause data spills, leading to write amplification.

  • During the shuffle read process, a large number of small network packets can cause connection reset errors.

  • The shuffle read process generates numerous small I/O requests and random reads, placing a high load on disks and the CPU.

  • When the number of mappers (M) and reducers (N) reaches the thousands, the M × N connection count makes job completion nearly impossible.

  • The NodeManager and the Spark External Shuffle Service run in the same process. When the shuffle data volume is extremely large, the NodeManager often restarts, which destabilizes YARN scheduling.

Celeborn addresses these shuffle process issues and provides the following benefits:

  • Uses a push-style shuffle instead of a pull-style shuffle to reduce memory pressure on mappers.

  • Supports I/O aggregation, which reduces the number of shuffle read connections from M × N to N and converts random reads into sequential reads.

  • Supports a two-replica mechanism to reduce the probability of fetch failures.

  • Supports a compute-storage separation architecture, allowing the shuffle service to be deployed in a dedicated hardware environment and isolated from the compute cluster.

  • Eliminates the dependency on local disk when running Spark on Kubernetes.

The following figure shows the Celeborn architecture.Celeborn

Prerequisites

You have created an EMR DataLake cluster or a custom cluster and selected the Celeborn service. For more information about how to create a cluster, see Create a cluster.

Limitations

This topic applies only to clusters of the following versions.

Cluster

Version

DataLake cluster

EMR-3.45.0 or later, and EMR-5.11.0 or later.

custom cluster

EMR-3.45.0 or later, and EMR-5.11.0 or later.

Procedure

Spark configuration

Parameter

Description

spark.shuffle.manager

  • For versions 0.4.x and later, the value must be org.apache.spark.shuffle.celeborn.SparkShuffleManager.

  • For versions 0.3.x and earlier, the value must be org.apache.spark.shuffle.celeborn.RssShuffleManager.

spark.serializer

The value must be org.apache.spark.serializer.KryoSerializer.

spark.celeborn.push.replicate.enabled

Specifies whether to enable the two-replica mechanism. Valid values:

  • true (default): Enables the two-replica mechanism.

  • false: Disables the two-replica mechanism.

spark.shuffle.service.enabled

Set this parameter to false to use Celeborn.

To use Celeborn, you must disable the existing External Shuffle Service. Spark's Dynamic Allocation feature works as expected when Celeborn is enabled.

Note
  • If spark.shuffle.service.enabled is set to true, Celeborn is not used.

  • Alibaba Cloud Spark and open source Spark 3.5 are compatible with Celeborn.

spark.celeborn.shuffle.writer

Celeborn supports the following writer modes:

  • hash (default): Uses more memory if the partition concurrency is high.

  • sort: Uses a fixed amount of memory and works stably even with high partition concurrency.

spark.celeborn.master.endpoints

Specify the endpoints in the format <celeborn-master-ip>:<celeborn-master-port>.

Parameters:

  • <celeborn-master-ip>: The public IP address of the master node.

  • <celeborn-master-port>: This value is fixed at 9097.

For a high-availability cluster, configure the IP addresses of all master nodes.

spark.sql.adaptive.enabled

Celeborn supports Adaptive Query Execution (AQE). For optimal shuffle performance, disable the local shuffle reader.

Set these parameters to true, false, and true, respectively.

spark.sql.adaptive.localShuffleReader.enabled

spark.sql.adaptive.skewJoin.enabled

The Spark service supports one-click configuration to use the Celeborn service.

  • For EMR-5.11.1 and later, and EMR-3.45.1 and later:

    On the Status page of the Spark service, in the Service Overview section, you can toggle the enableCeleborn switch.

  • For EMR-5.11.0 and EMR-3.45.0:

    On the Status page of the Spark service, in the Components section, find SparkThriftServer. In the Actions column, choose more > enableCeleborn or more > disableCeleborn. This action automatically modifies the Spark configuration parameters listed above, restarts SparkThriftServer, and updates both the spark-defaults.conf and spark-thriftserver.conf files.

    • If you select more > enableCeleborn, all Spark jobs use the Celeborn service.

    • If you select more > disableCeleborn, no Spark jobs use the Celeborn service.

Celeborn configuration

View and modify all Celeborn configuration parameters on the Celeborn service configuration page.

Important

Parameter values vary by node group (for example, CORE or TASK).

Parameter

Description

Default

celeborn.worker.flusher.threads

The number of threads for flushing data to disk (HDD or SSD).

  • HDD: 1

  • SSD: 8

CELEBORN_WORKER_OFFHEAP_MEMORY

The size of worker off-heap memory.

Automatically calculated based on the cluster configuration.

celeborn.application.heartbeat.timeout

The application heartbeat timeout. If the timeout is reached, the system releases the application's resources.

120 s

celeborn.worker.flusher.buffer.size

The flush buffer size. Flushing is triggered when this size is exceeded.

256 KB

celeborn.metrics.enabled

Specifies whether to enable monitoring. Valid values:

  • true: Enables monitoring.

  • false: Disables monitoring.

true

CELEBORN_WORKER_MEMORY

The size of worker heap memory.

1 GB

CELEBORN_MASTER_MEMORY

The size of master heap memory.

2 GB

Restart Celeborn components

  1. On the Status page of the Celeborn service, in the Actions column for the CelebornMaster component, select more > restart_clean_meta.

    Note

    For a non-high-availability cluster, you can also click Restart in the Actions column for the CelebornMaster component.

  2. In the dialog box, turn off the Rolling Execution switch, enter an execution reason, and click OK.

  3. In the confirmation dialog box, click OK.