This is an automated email from the ASF dual-hosted git repository.
dongjoon-hyun pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/spark.git
The following commit(s) were added to refs/heads/master by this push:
new 5300da3904d9 [SPARK-57503][CORE] Simplify removePendingChunks to avoid
redundant re-iteration in ShuffleBlockFetcherIterator
5300da3904d9 is described below
commit 5300da3904d936ac92817318df48a24b93fa6773
Author: YangJie <[email protected]>
AuthorDate: Wed Jun 17 09:32:47 2026 -0700
[SPARK-57503][CORE] Simplify removePendingChunks to avoid redundant
re-iteration in ShuffleBlockFetcherIterator
### What changes were proposed in this pull request?
`ShuffleBlockFetcherIterator.removePendingChunks` collected the removed
chunk ids with a `foreach { _ => ... }` whose body re-`flatMap`ped the whole
`fetchRequestsToRemove` queue on every iteration. With N removed requests it
repeated the same work N times; the result was still correct only because the
accumulator is a `HashSet` that absorbed the duplicates.
This processes each removed request once and drops the now-unnecessary
intermediate queue, so the helper becomes a single `dequeueAll(pred).foreach`
pass.
### Why are the changes needed?
To remove accidental O(N^2) work and a redundant intermediate collection on
the push-based shuffle fetch-failure fallback path, and to make the intent
clearer.
### Does this PR introduce _any_ user-facing change?
No. The returned set of removed chunk ids is unchanged.
### How was this patch tested?
Existing `ShuffleBlockFetcherIteratorSuite`, which covers the SPARK-32922
chunk-fetch fallback paths, still passes.
### Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 4.8)
Closes #56561 from LuciferYang/core-removePendingChunks-fix.
Authored-by: YangJie <[email protected]>
Signed-off-by: Dongjoon Hyun <[email protected]>
---
.../org/apache/spark/storage/ShuffleBlockFetcherIterator.scala | 9 +++------
1 file changed, 3 insertions(+), 6 deletions(-)
diff --git
a/core/src/main/scala/org/apache/spark/storage/ShuffleBlockFetcherIterator.scala
b/core/src/main/scala/org/apache/spark/storage/ShuffleBlockFetcherIterator.scala
index a2717ef39374..cb15e954bb38 100644
---
a/core/src/main/scala/org/apache/spark/storage/ShuffleBlockFetcherIterator.scala
+++
b/core/src/main/scala/org/apache/spark/storage/ShuffleBlockFetcherIterator.scala
@@ -1325,15 +1325,12 @@ final class ShuffleBlockFetcherIterator(
}
def filterRequests(queue: mutable.Queue[FetchRequest]): Unit = {
- val fetchRequestsToRemove = new mutable.Queue[FetchRequest]()
- fetchRequestsToRemove ++= queue.dequeueAll { req =>
+ queue.dequeueAll { req =>
val firstBlock = req.blocks.head
firstBlock.blockId.isShuffleChunk && req.address.equals(address) &&
sameShuffleReducePartition(firstBlock.blockId)
- }
- fetchRequestsToRemove.foreach { _ =>
- removedChunkIds ++=
-
fetchRequestsToRemove.flatMap(_.blocks.map(_.blockId.asInstanceOf[ShuffleBlockChunkId]))
+ }.foreach { req =>
+ removedChunkIds ++=
req.blocks.map(_.blockId.asInstanceOf[ShuffleBlockChunkId])
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]