Remote Shuffle with Celeborn#
See the Celeborn guide for setup and client compatibility. The settings below apply only to Comet’s native remote shuffle, which is unavailable with the currently released Celeborn 0.6.x and 0.7.x clients. They do not tune Celeborn’s existing Spark shuffle implementation.
Setting |
Default |
Purpose |
|---|---|---|
|
64 MiB |
Maximum size of one complete encoded frame. The admission budget can reduce the effective limit. |
|
512 MiB |
Shared memory admission budget for map attempts using the same executor-side remote shuffle client. |
Admission includes Arrow encoding workspace and overlapping native, JNI, and client frame copies. An ordinary uncompressed frame needs roughly seven times its size plus schema and codec overhead. The default 512 MiB budget accommodates ordinary frames up to the default 64 MiB frame limit. Compression reduces transmitted bytes but still needs workspace for the uncompressed data. This budget bounds shuffle-write admission, not total executor memory.
Comet splits batches between rows. If a row, schema, or encoding workspace cannot fit, Comet
replaces the exchange with its local native shuffle writer before downstream tasks consume
the output.
Increase the limiting setting if larger rows need to stay on the remote path, allowing for
encoding workspace as well as the encoded frame. Raising only maxFrameBytes may not help
when maxInFlightBytes is the limiting budget. Larger budgets also increase potential
executor memory use.
Native frames use Comet’s shuffle compression settings. The raw Celeborn path bypasses Celeborn’s additional row compression and decompression, so frames are not compressed twice.