You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 02:25:17