PySpark如何将多列合并为包含列表的单列?
PySpark合并多列为列表列的解决方案
你可以直接使用PySpark内置的array()函数来实现需求,这比自定义lambda/UDF更高效,且符合Spark的分布式计算优化逻辑。
步骤1:构造示例DataFrame
先创建和你描述一致的测试数据:
from pyspark.sql import SparkSession from pyspark.sql.functions import array # 初始化SparkSession spark = SparkSession.builder.appName("MergeColumns").getOrCreate() # 创建初始DataFrame data = [(11, 'a', 13), (21, 'b', 23)] df = spark.createDataFrame(data, ["Col1", "Col2", "Col3"]) df.show()
输出:
+----+----+----+ |Col1|Col2|Col3| +----+----+----+ | 11| a| 13| | 21| b| 23| +----+----+----+
步骤2:使用array()函数合并列
直接调用array()函数将Col2和Col3合并为列表列,再保留Col1即可:
result_df = df.select("Col1", array("Col2", "Col3").alias("NewCol2")) result_df.show(truncate=False)
输出:
+----+---------+ |Col1|NewCol2 | +----+---------+ |11 |[a, 13] | |21 |[b, 23] | +----+---------+
可选:自定义UDF实现(不推荐)
如果一定要用lambda风格的自定义函数,也可以通过udf()实现,但这种方式性能不如内置函数,因为UDF会打破Spark的Catalyst优化:
from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, StringType # 统一元素类型为字符串(自动转换整数为字符串) merge_udf = udf(lambda x, y: [x, y], ArrayType(StringType())) result_df_udf = df.select("Col1", merge_udf("Col2", "Col3").alias("NewCol2")) result_df_udf.show(truncate=False)
总结:优先使用Spark内置的array()函数,代码简洁且性能更优,完全不需要复杂的lambda操作。
内容的提问来源于stack exchange,提问作者Raphvanns
相关产品推荐
相关产品推荐

