在PySpark中合并多列至含索引列的单列
PySpark 多列合并并提取唯一有效值方案
实现步骤
- 导入依赖模块
from pyspark.sql import SparkSession from pyspark.sql.functions import array, array_remove, array_distinct, concat_ws, col
- 创建示例数据集
spark = SparkSession.builder.appName("MergeColumns").getOrCreate() data = [ (1, "A", None, None), (2, "B", "B", None), (3, None, "C", None), (4, None, None, "D") ] df = spark.createDataFrame(data, ["Index", "Column_1", "Column_2", "Column_3"])
- 核心处理逻辑
通过数组操作完成空值移除、重复值去重,最终合并为单列:
columns_to_merge = ["Column_1", "Column_2", "Column_3"] result_df = df.withColumn( "Value", concat_ws("", array_distinct(array_remove(array(*[col(c) for c in columns_to_merge]), None))) ).select("Index", "Value") result_df.show()
关键函数说明
array(*[col(c) for c in columns_to_merge]):将目标列转换为数组格式array_remove(..., None):过滤数组中的所有空值array_distinct(...):对数组内的元素去重,解决同一行多列重复值的问题concat_ws("", ...):将去重后的数组元素拼接为单个字符串(处理后每行数组仅含一个有效值)
内容的提问来源于stack exchange,提问作者Mikesama
相关产品推荐
相关产品推荐

