[SPARK-58911][ML] Use DataFrame aggregations in CountVectorizer fit - #57755
Closed
zhengruifeng wants to merge 6 commits into
Closed
[SPARK-58911][ML] Use DataFrame aggregations in CountVectorizer fit#57755zhengruifeng wants to merge 6 commits into
zhengruifeng wants to merge 6 commits into
Conversation
zhengruifeng
marked this pull request as ready for review
August 20, 2026 15:02
HyukjinKwon
approved these changes
Aug 20, 2026
zhengruifeng
added a commit
that referenced
this pull request
Aug 21, 2026
### What changes were proposed in this pull request? This PR rewrites `CountVectorizer.fit` to stay in the DataFrame execution path instead of converting the input column to an RDD. The new implementation: - computes fractional `minDF` / `maxDF` thresholds with a scalar subquery over the input size; - computes word counts with `explode` and grouped aggregation; - computes document frequency, when DF filtering is requested, with a two-stage aggregation: first by `(docId, word)` and then by `word`; - uses DataFrame ordering and `limit` to select the vocabulary. ### Why are the changes needed? The current implementation builds per-document maps manually and combines them with `reduceByKey`. Keeping the fit logic in DataFrame operations lets Spark plan the aggregations natively while preserving the existing CountVectorizer semantics for word counts, document frequency filtering, and deterministic vocabulary ordering. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Added regression coverage for document frequency when a document contains repeated words. ``` ./build/sbt "mllib/testOnly org.apache.spark.ml.feature.CountVectorizerSuite" ``` ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Codex (GPT-5) Closes #57755 from zhengruifeng/count-vectorizer-dataframe-fit-dev3. Authored-by: Ruifeng Zheng <ruifengz@apache.org> Signed-off-by: Ruifeng Zheng <ruifengz@foxmail.com> (cherry picked from commit 619c067) Signed-off-by: Ruifeng Zheng <ruifengz@foxmail.com>
Contributor
Author
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?
This PR rewrites
CountVectorizer.fitto stay in the DataFrame execution path instead ofconverting the input column to an RDD.
The new implementation:
minDF/maxDFthresholds with a scalar subquery over the input size;explodeand grouped aggregation;first by
(docId, word)and then byword;limitto select the vocabulary.Why are the changes needed?
The current implementation builds per-document maps manually and combines them with
reduceByKey.Keeping the fit logic in DataFrame operations lets Spark plan the aggregations natively while
preserving the existing CountVectorizer semantics for word counts, document frequency filtering,
and deterministic vocabulary ordering.
Does this PR introduce any user-facing change?
No.
How was this patch tested?
Added regression coverage for document frequency when a document contains repeated words.
Was this patch authored or co-authored using generative AI tooling?
Generated-by: Codex (GPT-5)