Replace deprecated stateful UDAFs with typed aggregators - #760
Merged
Conversation
kyraman
approved these changes
Jul 17, 2026
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.
Issue #, if available:
Fixes #583
Description of changes:
This replaces the deprecated
UserDefinedAggregateFunctionimplementations for theDataTypeandKLLSketchanalyzers with Spark's typedAggregatorAPI.The main improvement is in KLL processing. The old implementation serialized and deserialized the sketch for every input row. The new aggregator keeps the sketch in memory while processing a partition and only serializes it when Spark needs to encode the aggregation buffer.
DataTypenow uses a typed buffer containing the same five counters as before.The existing output formats and behavior are unchanged. I added tests that compare the old and new implementations byte for byte, including null and empty inputs, shuffled partitions, KLL compaction, special floating-point values, separate physical plans, and object-hash fallback.
Validation performed:
StatefulAggregatorsTest: 7 tests passedAnalyzerTests: 98 tests passedKLLProfileTest: 3 tests passedIn a small local Spark 3.5.7 benchmark, KLL was about 11% faster.
DataTypewas about 12% slower, so this change should not be considered a performance improvement for that analyzer; it does remove its use of the deprecated API.By submitting this pull request, I confirm that my contribution is made under the terms of the Apache 2.0 license.