基于account_type字段使用PySpark合并多行数据为单行
PySpark实现基于account_type的列合并方案
核心思路
要实现需求,核心是按Id聚合不同account_type下的字段,再将指定类型的字段拆分为独立列,同时处理空值场景:
- 以
account_type=1的数据作为主列(col_a、col_b) - 将
account_type=2及其他类型的数据映射为新扩展列(如col_a_account_type_2) - 无对应类型数据时,用空字符串填充
代码实现
1. 初始化环境与测试数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 初始化Spark会话 spark = SparkSession.builder.appName("account_merge").getOrCreate() # 定义数据Schema schema = StructType([ StructField("Id", StringType(), nullable=True), StructField("col_a", StringType(), nullable=True), StructField("col_b", StringType(), nullable=True), StructField("account_type", IntegerType(), nullable=True) ]) # 构造测试数据集 data = [ ("James Butt", "a1_col_a_data_1", "a1_col_b_data_1", 1), ("James Butt", "a1_col_a_data_2", "a1_col_b_data_2", 2), ("Art Venere", "a1_col_a_data_3", "a1_col_b_data_3", 1), ("Lenna Paprocki", "a1_col_a_data_4", "a1_col_b_data_4", 1), ("Lenna Paprocki", "a1_col_a_data_5", "a1_col_b_data_5", 2), ("John Doe", "a1_col_a_data_6", "a1_col_b_data_6", 1), ("Mitsue Tollner", "a1_col_a_data_7", "a1_col_b_data_7", 1), ("Leota Dilliard", "a1_col_a_data_8", "a1_col_b_data_8", 1), ("Sage Wieser", "a1_col_a_data_9", "a1_col_b_data_9", 1), ("Sage Wieser", "a1_col_a_data_10", "a1_col_b_data_10", 2) ] df = spark.createDataFrame(data, schema=schema)
2. 基础实现(针对account_type=2的扩展)
# 按Id分组,收集所有account_type对应的字段结构体 grouped_df = df.groupBy("Id").agg( F.collect_list(F.struct("account_type", "col_a", "col_b")).alias("account_data") ) # 提取基础数据与扩展数据,并映射为目标列 result_df = grouped_df \ .withColumn("base_data", F.filter("account_data", lambda x: x.account_type == 1).getItem(0)) \ .withColumn("ext_data", F.filter("account_data", lambda x: x.account_type == 2).getItem(0)) \ .select( "Id", F.col("base_data.col_a").alias("col_a"), F.col("base_data.col_b").alias("col_b"), # 空值替换为空白字符串 F.coalesce(F.col("ext_data.col_a"), F.lit("")).alias("col_a_account_type_2"), F.coalesce(F.col("ext_data.col_b"), F.lit("")).alias("col_b_account_type_2") ) # 查看结果 result_df.show(truncate=False)
3. 扩展支持最多10种account_type
如果需要支持1-10所有account_type(以1为基础,2-10为扩展),可以通过循环自动生成列:
base_type = 1 ext_types = list(range(2, 11)) # 2到10的account_type # 分组收集数据 grouped_df = df.groupBy("Id").agg( F.collect_list(F.struct("account_type", "col_a", "col_b")).alias("account_data") ) # 初始化结果集,提取基础列 result_df = grouped_df \ .withColumn("base_data", F.filter("account_data", lambda x: x.account_type == base_type).getItem(0)) \ .select( "Id", F.col("base_data.col_a").alias("col_a"), F.col("base_data.col_b").alias("col_b") ) # 循环添加所有扩展类型的列 for t in ext_types: result_df = result_df \ .withColumn(f"temp_{t}", F.filter("account_data", lambda x: x.account_type == t).getItem(0)) \ .withColumn(f"col_a_account_type_{t}", F.coalesce(F.col(f"temp_{t}.col_a"), F.lit(""))) \ .withColumn(f"col_b_account_type_{t}", F.coalesce(F.col(f"temp_{t}.col_b"), F.lit(""))) \ .drop(f"temp_{t}") # 清理临时列 result_df = result_df.drop("account_data") result_df.show(truncate=False)
代码说明
- 分组收集:使用
groupBy+collect_list将同一Id下的所有account_type数据聚合为结构体列表,方便后续筛选。 - 数据筛选:通过
F.filter精准提取指定account_type的结构体数据。 - 空值处理:
coalesce函数确保无对应类型数据时,用空白字符串替代默认的null,匹配需求输出格式。
内容的提问来源于stack exchange,提问作者Eyedia
相关产品推荐
相关产品推荐

