PySpark按CLMN_SEQ_NUM排序,将多行CLMN_NM合并为单行逗号分隔串
按指定列排序后合并DataFrame列值的解决方案
问题描述
我有一个包含CLMN_SEQ_NUM和CLMN_NM两列的DataFrame,想要将CLMN_NM列的多行值合并为单行,以逗号分隔。期望输出为:PR_NAME,PR_ID,PR_ZIP,PR_ADDRESS,PR_COUNTRY。
使用如下代码后得到的结果顺序不符合预期:
cols_comb = df.agg(F.concat_ws(",",F.collect_list(F.col("CLMN_NM")))).first()[0]
输出结果为:PR_ZIP,PR_NAME,PR_COUNTRY,PR_ID,PR_ADDRESS,需要按CLMN_SEQ_NUM排序后再合并列值。
解决方案
Spark的collect_list函数本身不保证顺序,想要按CLMN_SEQ_NUM排序后合并,有两种可靠的实现方式:
方式一:利用结构体数组排序(推荐)
将序号和名称打包成结构体,先对结构体数组按序号排序,再提取名称列进行合并,这种方式不受分区影响,结果更稳定:
from pyspark.sql import functions as F cols_comb = df.agg( F.concat_ws( ",", F.transform( F.sort_array(F.collect_list(F.struct("CLMN_SEQ_NUM", "CLMN_NM"))), lambda x: x["CLMN_NM"] ) ) ).first()[0]
F.struct("CLMN_SEQ_NUM", "CLMN_NM"):将序号和名称组合成结构体,保留排序依据F.sort_array(...):按结构体中的CLMN_SEQ_NUM升序排序整个数组F.transform(...):从排序后的结构体数组中提取出CLMN_NM的值F.concat_ws(",", ...):将提取出的名称列表用逗号连接成目标字符串
方式二:先排序DataFrame再收集合并
如果你的DataFrame数据量较小,或可以保证分区不影响排序结果,也可以先对整个DataFrame按CLMN_SEQ_NUM排序,再执行合并:
cols_comb = df.orderBy("CLMN_SEQ_NUM").agg( F.concat_ws(",", F.collect_list("CLMN_NM")) ).first()[0]
若DataFrame存在多分区,建议先合并分区再排序,避免分区内排序导致整体顺序错误:
cols_comb = df.repartition(1).orderBy("CLMN_SEQ_NUM").agg( F.concat_ws(",", F.collect_list("CLMN_NM")) ).first()[0]
内容的提问来源于stack exchange,提问作者newbie
相关产品推荐
相关产品推荐

