如何在PySpark中将多列去重值合并为新列?
PySpark实现新增包含所有列去重值的统一列
实现步骤与代码
- 初始化SparkSession并创建示例数据集
from pyspark.sql import SparkSession from pyspark.sql.functions import lit, collect_set, array_distinct, flatten # 初始化SparkSession spark = SparkSession.builder.appName("AddUnifiedColumn").getOrCreate() # 构建示例DataFrame data = [("a", "e", "i"), ("b", "f", "j"), ("c", "g", "k"), ("d", "h", "l")] df = spark.createDataFrame(data, ["C1", "C2", "C3"])
- 生成所有列的去重值拼接字符串
通过Spark内置函数收集所有列的去重值,合并后再次去重,最终拼接成空格分隔的字符串:
# 收集各列去重值→合并数组→去重→转为空格分隔字符串 distinct_values_array = df.select( array_distinct(flatten([collect_set(col) for col in df.columns])) ).first()[0] c4_content = " ".join(distinct_values_array)
- 添加统一的C4列
使用lit函数将拼接好的字符串作为常量列添加到原DataFrame的每一行:
result_df = df.withColumn("C4", lit(c4_content)) # 查看结果 result_df.show(truncate=False)
关键函数说明
collect_set(col):提取单列的所有去重值,返回数组类型flatten():将多个列的数组合并为一维数组array_distinct():对合并后的数组做全局去重,避免跨列重复值lit():将常量值转为Spark列,实现每行添加相同内容
运行后输出结果:
+---+---+---+------------------------+ |C1 |C2 |C3 |C4 | +---+---+---+------------------------+ |a |e |i |a b c d e f g h i j k l | |b |f |j |a b c d e f g h i j k l | |c |g |k |a b c d e f g h i j k l | |d |h |l |a b c d e f g h i j k l | +---+---+---+------------------------+
内容的提问来源于stack exchange,提问作者Nabs335
相关产品推荐
相关产品推荐

