Skip to content

Shuffle tracking prevents dynamic allocation from removing executors after collect() / toPandas() #59014

Description

@aalopatin

When using dynamic allocation with shuffle tracking enabled, executors that participated in a shuffle remain alive indefinitely after a collect() or toPandas() action has completed.

Configuration:

spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.executorIdleTimeout=60s
spark.dynamicAllocation.cachedExecutorIdleTimeout=1m
spark.dynamicAllocation.shuffleTracking.enabled=true
spark.dynamicAllocation.shuffleTracking.timeout=infinity

For example:

pandas_df = spark.sql("...").toPandas()

After the action completes, the executors that produced shuffle data are not removed after spark.dynamicAllocation.executorIdleTimeout.

The same behavior can be reproduced with:

df.collect()

I also observed that shuffle files are present in $SPARK_LOCAL_DIRS on the executor pods (shuffle_*.data, shuffle_*.index, etc.).

Interestingly, setting a finite value for:

spark.dynamicAllocation.shuffleTracking.timeout

allows the executors to be removed, which suggests that shuffle tracking is what prevents their removal.

My question is whether this behavior is expected.

collect() and toPandas() are terminal actions that materialize their result on the driver. Once the action has completed successfully, there should presumably be no further consumer of the shuffle output produced by the query.

Therefore, I would expect the executors to become eligible for removal after spark.dynamicAllocation.executorIdleTimeout, even if shuffle files produced during the query still physically exist on the executors.

Is there a reason why the shuffle is still considered active after the collect() / toPandas() action has completed? Or is this an expected limitation/behavior of shuffle tracking?

For comparison, setting:

spark.dynamicAllocation.shuffleTracking.timeout=1m

allows the executors to be removed, whereas:

spark.dynamicAllocation.shuffleTracking.timeout=infinity

keeps them alive.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions