PySpark按streetID分组生成排序houseNum列表的实现方案咨询
PySpark分组聚合生成排序后列表的解决方案
你要的需求完全不需要使用自定义UDF,PySpark提供了现成的内置API可以高效实现,性能远高于自定义Python UDF,避免了不必要的序列化/反序列化开销。
实现代码
首先导入依赖的内置函数:
from pyspark.sql import functions as F
直接对DataFrame做分组聚合即可:
# 假设你的原始DataFrame名称为df result_df = df.groupBy("streetID") \ .agg( F.array_sort(F.collect_list("houseNum")).alias("houseNumList") ) # 如果需要和示例一致按streetID升序排列结果,可补充orderBy result_df = result_df.orderBy("streetID", ascending=True)
函数说明
groupBy("streetID"):按streetID字段对数据做分组,相同streetID的行会被划分到同一处理组collect_list("houseNum"):将同一分组内所有houseNum的值收集为一个列表,默认列表顺序与数据读取/处理的行顺序一致array_sort():对传入的列表做升序排序,直接得到你需要的排序后列表结果alias("houseNumList"):将聚合后的结果列重命名为目标列名
补充说明
array_sort为Spark 2.4及以上版本提供的内置函数,如果你使用的是更早的Spark版本,才需要额外自定义UDF实现排序逻辑,优先建议使用内置函数,性能比Python UDF高5~10倍。- 如果你的业务场景中同一streetID下存在重复的houseNum需要去重,可以将
collect_list替换为collect_set,后续再搭配排序即可得到去重后的有序列表。
内容的提问来源于stack exchange,提问作者Tony LaRussa
相关产品推荐
相关产品推荐

