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:
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.
When using dynamic allocation with shuffle tracking enabled, executors that participated in a shuffle remain alive indefinitely after a
collect()ortoPandas()action has completed.Configuration:
For example:
After the action completes, the executors that produced shuffle data are not removed after
spark.dynamicAllocation.executorIdleTimeout.The same behavior can be reproduced with:
I also observed that shuffle files are present in
$SPARK_LOCAL_DIRSon the executor pods (shuffle_*.data,shuffle_*.index, etc.).Interestingly, setting a finite value for:
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()andtoPandas()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:
allows the executors to be removed, whereas:
keeps them alive.