[SPARK-58939][PYSPARK] Move Arrow collect batch reordering into ArrowCollectSerializer - #58221
Open
Yicong-Huang wants to merge 1 commit into
Open
[SPARK-58939][PYSPARK] Move Arrow collect batch reordering into ArrowCollectSerializer#58221Yicong-Huang wants to merge 1 commit into
Yicong-Huang wants to merge 1 commit into
Conversation
uros-b
approved these changes
Aug 22, 2026
Member
|
LGTM, thank you @Yicong-Huang! |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
Move the batch-reordering logic for Arrow-based collect from the
_collect_as_arrowcaller intoArrowCollectSerializer.Dataset.collectAsArrowToPython(Spark classic) streams Arrow batches to Python out of order and appends batch-order indices at the end. Previously the serializer yielded raw batches plus the order list, and_collect_as_arrowsplit off the order list and reordered. NowArrowCollectSerializerextendsArrowStreamSerializerand yields the batches already ordered, so the caller drops the results slicing, the reorder comprehension, and theisinstancecheck. The unuseddump_streamdelegate is also removed. Net -14 lines.To keep
spark.sql.execution.arrow.pyspark.selfDestruct.enabledeffective, the serializer drops its reference to each batch as it yields it (batches[i] = None; the indices are a permutation), so peak memory stays ~1x.Why are the changes needed?
The ordering protocol is an implementation detail of the collect stream; owning it inside the serializer removes duplicated logic from
_collect_as_arrowand drops dead code, with no behavior change.Does this PR introduce any user-facing change?
No.
How was this patch tested?
Existing
python/pyspark/sql/tests/arrow/test_arrow.py(toPandas/toArrow + self-destruct) in CI; also validated reorder, JVM-error, empty-stream, and repr against a simulated stream.Was this patch authored or co-authored using generative AI tooling?
No.