Hi Zaynt,

Thanks for digging up CELEBORN-1896 / PR #3109 and the 2088 follow-up.
I'll open a JIRA and PR to flip the default, targeting master with a backport 
request to branch-0.7. 

Best,
Randy

> 2026年8月3日 11:30,Zhentao Shuai <[email protected]> 写道:
> 
> Hi Randy,
> Thanks for raising this — the production numbers make a compelling case. I 
> read through the implementation to ground my view.
> The feature was introduced by [CELEBORN-1896] (PR #3109, merged into 0.6.0) 
> and is present on branch-0.7, including the follow-up NPE fix [CELEBORN-2088] 
> (PR #3401).
> Why false by default. It was a new feature in 0.6.0 operating on the 
> data-deletion path, where a wrong deletion means data loss, not a no-op. 
> Opt-in is the standard conservative choice for first-cut safety-sensitive 
> cleanup; it signals "unproven," not "unsafe."
> I support flipping to true. stageRerun.enabled has defaulted to true since 
> 0.4.0, so retries are already the out-of-the-box behavior, yet 
> cleanFailedShuffle=false means stale data from each retried attempt is never 
> reclaimed promptly. The two defaults contradict each other, and the failure 
> mode you described (300+ TiB of stale data, disk-full cascades) is exactly 
> what this feature removes. The mechanism reuses the existing 
> UnregisterShuffle path with no new deletion route.
> Welcome to file a JIRA / PR to flip this config.
> Best,
> Zaynt Shuai
> 
> On 2026/07/30 07:38:53 Hua Hai wrote:
>> Hi community,
>> 
>> I'd like to start a discussion around the default value of
>> `celeborn.client.spark.fetch.cleanFailedShuffle`, introduced in 0.6.0.
>> 
>> Background
>> ----------
>> 
>> The config gates a cleanup path that reclaims disk space held by
>> shuffle files which can no longer be fetched. The mechanism itself
>> looks safe:
>> 
>>  * When a FetchFailed is captured on the driver, the failed shuffleId
>>    is recorded.
>>  * When the retried stage registers a new shuffleId, the previously
>>    failed shuffleId is enqueued for cleanup.
>>  * FailedShuffleCleaner then reuses the existing UnregisterShuffle
>>    path to remove the stale shuffle data on the workers.
>> 
>> Since 0.4.0, `celeborn.client.spark.stageRerun.enabled` defaults to
>> `true`, so FetchFailed stage rerun is already the out-of-the-box
>> behavior. However `cleanFailedShuffle` still defaults to `false`,
>> which means every retried stage silently leaves its old shuffle data
>> on the workers until the app exits or the periodic expired-shuffle
>> cleanup catches it.
>> 
>> Why this matters in production
>> ------------------------------
>> 
>> We've hit the disk-pressure side of this on large shuffle jobs where
>> a stage retries multiple times. A recent case from our cluster:
>> 
>>  * A single stage retried three times.
>>  * Retry #1 wrote ~143 TiB of shuffle data before failing.
>>  * Retry #2 wrote another ~143 TiB before failing.
>>  * Retry #3 was still running, ~61 TiB written so far.
>> 
>> With `cleanFailedShuffle=false`, all three generations of shuffle data
>> sit on the workers simultaneously — well over 300 TiB of stale bytes
>> that the job will never fetch again. On a shared cluster this drives
>> workers into disk-full / high-usage exclusion, cascading into more
>> FetchFailed for unrelated apps and destabilizing the cluster until
>> the app finally exits. Turning `cleanFailedShuffle` on effectively
>> removes this failure mode for us.
>> 
>> Questions for the community
>> ---------------------------
>> 
>>  1. What was the rationale for keeping `cleanFailedShuffle` disabled
>>     by default when it was introduced?
>> 
>>  2. Would the community be open to flipping the default to `true`
>>     in a future release?
>> 
>>  3. Are there scenarios where enabling this flag could cause a
>>     late-arriving fetch to hit a "shuffle not registered" error
>>     instead of the current benign leak?
>> 
>> Happy to follow up with a JIRA / doc PR / default-flip PR 
>> once the direction is clear.
>> 
>> Thanks,
>> Haihua Wang (Randy)

Reply via email to