You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark实现同Service行合并:生成带Usage标签的Info列

PySpark DataFrame 按Service聚合生成结构化Usage计数字段

问题背景

原始PySpark DataFrame结构如下(数据顺序不固定,Service数量无上限):

serviceusagecount
alow5
ahigh3
bhigh3
blow2

需要生成新DataFrame,每行包含对应Service的Usage计数结构化字符串,目标输出如下:

serviceinfo
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 22:22:27