Pyspark将同一客户多条更新DF聚合为单条JSON结构问题
PySpark 单客户多变更记录合并为单条JSON解决方案
问题根因
你当前输出多条JSON的核心原因是groupBy中加入了PhonePreference和EmailPreference字段,三行原始记录中这两个字段的组合存在两种不同值(<null, Primary>、<Primary, null>),所以会拆分出多个分组,生成多条结果。
调整思路
- 仅按
CustomerId作为分组维度,保证每个客户仅生成一个分组 - 构造结构体时使用
when条件过滤:只有手机号相关变更行才生成合法的PhoneDetails结构体,邮箱变更行对应的PhoneDetails设为null;同理处理EmailDetails - 聚合收集列表时过滤掉null值,即可将同类型所有变更归入对应列表
调整后代码
from pyspark.sql import functions as F json_df = updatesDF.select( F.col("CustomerId"), # 条件生成PhoneDetails:仅属性为手机号相关变更时生成,否则为null F.when( F.col("attribute").isin("PhoneNumber", "FormatPhoneNumber"), F.struct( F.col("PhoneType"), F.col("PhonePreference"), F.col("attribute"), F.col("from"), F.col("to") ) ).alias("PhoneDetails"), # 条件生成EmailDetails:仅属性为邮箱相关变更时生成,否则为null F.when( F.col("attribute") == "EmailAddress", F.struct( F.col("EmailType"), F.col("EmailPreference"), F.col("attribute"), F.col("from"), F.col("to") ) ).alias("EmailDetails") ).groupBy("CustomerId").agg( # 收集列表时过滤空值,只保留有效变更记录 F.collect_list(F.expr("filter(PhoneDetails, x -> x is not null)")).alias("PhoneDetails"), F.collect_list(F.expr("filter(EmailDetails, x -> x is not null)")).alias("EmailDetails") ) # 导出为JSON即可得到每个客户单条的结果 # json_df.write.json("输出路径")
输出效果
调整后每个客户仅生成1条JSON,你示例的C1000001客户最终JSON结构如下:
{ "CustomerId": "C1000001", "PhoneDetails": [ {"PhoneType":"Home","PhonePreference":"Primary","attribute":"PhoneNumber","from":"8177777777","to":"8168888888"}, {"PhoneType":"Home","PhonePreference":"Primary","attribute":"FormatPhoneNumber","from":"(816)777-7777","to":"(816)888-8888"} ], "EmailDetails": [ {"EmailType":"Home","EmailPreference":"Primary","attribute":"EmailAddress","from":"TEST@Solutions.com","to":"WELL@Solutions.com"} ] }
内容的提问来源于stack exchange,提问作者Sneha Nair
相关产品推荐
相关产品推荐

