Remote Shuffle Service (RSS) is an E-MapReduce extension that improves the stability and performance of the native Spark shuffle process. You can associate a Spark cluster on EMR on ACK with RSS to optimize shuffle operations.
Background information
Running Spark shuffle in a Container Service for Kubernetes (ACK) environment can cause the following issues:
-
Spark shuffle depends on local storage. Many compute-storage separated instance types and elastic container instance (ECI) scenarios lack built-in local disks, requiring you to purchase and attach cloud disks at additional cost, which reduces efficiency.
-
Spark 2 does not support dynamic allocation in ACK environments. Spark 3 implements dynamic allocation based on ShuffleTracking, but executor recycling is inefficient.
The native Spark shuffle implementation has the following drawbacks:
- A data overflow occurs if a large amount of data exists in a shuffle write task. This causes write amplification.
- A large number of small-size network packets exist in a shuffle read task. This causes connection reset.
- A large number of small-size I/O requests and random reads exist in a shuffle read task. This causes high disk and CPU loads.
- If thousands of mappers (M) and reducers (N) are used, a large number of connections are generated, which makes it difficult for jobs to run. The number of connections is calculated by using the following formula: M × N.
RSS from E-MapReduce resolves these shuffle issues and fully supports dynamic allocation in ACK environments.
Prerequisites
-
You have created a Spark cluster on the EMR on ACK page. For more information, see Step 1: Create a cluster.
-
You have created a Shuffle Service cluster on the EMR on ACK page. For more information, see Step 1: Create a cluster.
Limitations
-
A Spark cluster can only be associated with a Shuffle Service cluster in the same ACK cluster.
-
When you associate a Spark cluster on EMR on ACK with an RSS cluster on EMR on ACK, ensure they have the same cluster version to avoid compatibility risks. You can view the cluster version on the Cluster Details page.
Procedure
-
Log on to the EMR on ACK page.
-
Associate the Spark cluster with RSS.
-
On the EMR on ACK page, click the name of the Spark cluster.
-
On the Cluster Details page, in the Basic Information section, click Associate Now next to Associate RSS Cluster.
-
In the Associated Cluster section, click Add.
-
In the Associated Cluster dialog box, select the Shuffle Service cluster that you created, and then click Associate.
-
-
Optional: Configure RSS parameters. For more information, see Celeborn configurations.