Spark UDAF能否用于Window?测试几何平均UDAF是否合规?
Great question! Let’s unpack this for you clearly:
Absolutely. While it’s easy to miss because official Spark docs don’t call this out prominently, custom UDAFs are fully supported as window functions in Spark 2.x (including 2.1).
Window functions in Spark are fundamentally built on aggregation logic—they compute values over a sliding or fixed window of rows for each record. As long as your UDAF properly implements the UserDefinedAggregateFunction interface (with correct initialize, update, merge, and evaluate methods, plus properly defined input/buffer/output data types), it can work exactly like built-in window aggregations (e.g., sum() OVER (...), avg() OVER (...)).
If your geometric mean UDAF works correctly in standard aggregation scenarios, its output as a window function should be reliable too—provided your UDAF implementation handles window-specific logic properly:
- Ensure the
updatemethod correctly processes each row in the window - Verify the
mergemethod properly combines buffer data when windows span Spark partitions - Confirm the
evaluatemethod correctly calculates the final geometric mean from the buffer state
A quick way to validate: use a small, controlled dataset where you can manually compute the geometric mean for each window, then compare those results to what your UDAF returns.
Even though you didn’t finish this part of your question, a common point of confusion here is how UDAFs behave across these two scenarios. The core aggregation logic stays the same:
- In standard
GROUP BYaggregation, the UDAF computes a single value per group - In window aggregation, the UDAF computes a value for each row, based on the rows in its assigned window
As long as your UDAF’s core logic is sound, both use cases should produce accurate results.
内容的提问来源于stack exchange,提问作者Raphael Roth

