如何在PySpark DataFrame中按分组数据新增存储兴趣列表的列
PySpark 分组拼接多值为分隔符字符串实现方案
核心逻辑是按name字段分组后,先将同组的interest值收集为数组,再通过分号拼接为单字符串,可直接运行以下完整示例:
依赖导入
from pyspark.sql import SparkSession from pyspark.sql.functions import collect_list, concat_ws, collect_set
完整实现代码
# 初始化Spark会话 spark = SparkSession.builder.appName("merge_interests").getOrCreate() # 模拟示例数据,实际使用时替换为自己的DataFrame即可 raw_data = [("A", "gym"), ("A", "food"), ("A", "games"), ("B", "games")] df = spark.createDataFrame(raw_data, schema=["name", "interest"]) # 核心处理逻辑 result_df = df.groupBy("name") \ .agg( concat_ws(";", collect_list("interest")).alias("interests") ) # 打印结果验证 result_df.show(truncate=False)
说明
collect_list("interest"):按分组收集所有兴趣值,保留原始顺序和重复值,符合示例的输出要求- 如果需要对同用户的重复兴趣去重,将
collect_list替换为collect_set即可 concat_ws(";", 数组):用指定的分号作为分隔符,将数组元素拼接为单个字符串,自动忽略null值,无需额外处理空兴趣记录
运行后输出结果和要求的结构完全一致:
+----+---------------+ |name|interests | +----+---------------+ |A |gym;food;games | |B |games | +----+---------------+
内容的提问来源于stack exchange,提问作者timo
相关产品推荐
相关产品推荐

