How to enforce a specific partitioning and sort order automatically with a UDAF in Spark?

Viewed 26

I think the title of the question is quite clear. Based on the documentation regarding custom UDAFs, I want to develop a UDAF which uses a certain algorithm that relies on the fact that the reduce(b: BUF, a: IN): BUF function is called for inputs that are partitioned and sorted on the field that is being aggregated. And if this is not the case, I would like this to be enforced based on the fact that I'm using this UDAF, rather than repartitioning and sorting manually.

For example: If I were to develop my own my_count_distinct; say I'm calling it like this: my_count_distinct("user_id").

  • The reduce would increment a counter in the buffer when the current user_id is different than the previous one, as the algorithm assumes the the input is ordered on user_id.
  • The merge would just simply add up the counters in the buffer since the assumption is that the DataFrame is partitioned on the user_id, therefore the input per partition should be mutually exclusive.
0 Answers
Related