PySpark实现同Service行合并:生成带Usage标签的Info列
PySpark DataFrame 按Service聚合生成结构化Usage计数字段
问题背景
原始PySpark DataFrame结构如下(数据顺序不固定,Service数量无上限):
| service | usage | count |
|---|---|---|
| a | low | 5 |
| a | high | 3 |
| b | high | 3 |
| b | low | 2 |
需要生成新DataFrame,每行包含对应Service的Usage计数结构化字符串,目标输出如下:
| service | info |
|---|---|
| a | {low: 5, high: 3} |
| b | {low: 2, high: 3} |
尝试的代码仅收集了count的数值列表,无法对应到具体的usage标签:
# data is originally in dataframe called df: new_df = df.groupBy('service').agg(F.collect_list('count').alias('info'))
解决方案
方法1:生成Map类型后转字符串(推荐,支持后续结构化操作)
先将usage和count组合成键值对结构体,再转换为Map类型,最后转成目标字符串格式:
from pyspark.sql import functions as F # 生成Map类型字段,再转成无引号的目标格式 new_df = df.groupBy("service") \ .agg(F.map_from_entries(F.collect_list(F.struct("usage", "count"))).alias("info_map")) \ .withColumn("info", F.regexp_replace(F.to_json("info_map"), '"', '')) \ .drop("info_map")
这里用regexp_replace去掉了JSON字符串的引号,让输出和示例完全一致。如果不需要去掉引号,直接保留to_json的结果即可。
方法2:直接拼接字符串片段
先把每个usage和count拼接成low: 5这样的片段,再收集所有片段并用逗号分隔,最后包裹大括号:
from pyspark.sql import functions as F new_df = df.groupBy("service") \ .agg(F.collect_list(F.concat(F.col("usage"), F.lit(": "), F.col("count"))).alias("info_parts")) \ .withColumn("info", F.concat(F.lit("{"), F.concat_ws(", ", "info_parts"), F.lit("}"))) \ .drop("info_parts")
可选:固定输出顺序
如果需要强制low在前、high在后的顺序,可以在聚合前先排序:
from pyspark.sql import functions as F new_df = df.orderBy("service", "usage") \ .groupBy("service") \ .agg(F.collect_list(F.concat(F.col("usage"), F.lit(": "), F.col("count"))).alias("info_parts")) \ .withColumn("info", F.concat(F.lit("{"), F.concat_ws(", ", "info_parts"), F.lit("}"))) \ .drop("info_parts")
内容的提问来源于stack exchange,提问作者cliqer
相关产品推荐
相关产品推荐

