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

如何用Pyspark按ID合并DataFrame多行并取各字段最新更新值

PySpark 按ID合并多行取各字段最新有效值方案

需求说明

按ID分组合并多行数据,规则如下:

  • 字段值为NULL或者<*no-update*>时,视为本次未更新该字段
  • 每个字段取该ID下最新的有效值
  • 最终输出的updated_at取该ID下最大的时间戳

实现代码

环境准备与测试数据构造

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化SparkSession
spark = SparkSession.builder.appName("merge_rows_by_id").getOrCreate()

# 构造示例输入数据
data = [
    (123, "update1", "<*no-update*>", 1634228709),
    (123, "<*no-update*>", "80", 1634228724),
    (123, "update2", "<*no-update*>", 1634229000)
]
df = spark.createDataFrame(data, schema=["id", "column1", "column2", "updated_at"])

实现方式1:窗口函数(适合字段少的场景)

通过分区窗口按时间倒序排列,直接取每个字段第一个有效值:

# 定义窗口规则:按ID分组,更新时间倒序排列
w = Window.partitionBy("id").orderBy(F.desc("updated_at")) \
    .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)

result_df = df.withColumn("column1", 
        F.first(F.when(F.col("column1") != "<*no-update*>", F.col("column1")), ignorenulls=True).over(w)
    ) \
    .withColumn("column2", 
        F.first(F.when(F.col("column2") != "<*no-update*>", F.col("column2")), ignorenulls=True).over(w)
    ) \
    .filter(F.row_number().over(Window.partitionBy("id").orderBy(F.desc("updated_at"))) == 1) \
    .select("id", "column1", "column2", "updated_at")

result_df.show()

实现方式2:分组聚合(适合字段多的场景,可自动遍历处理)

无需手动编写每个字段的处理逻辑,自动遍历所有业务字段完成合并:

# 提取需要合并的业务字段(排除主键和时间字段)
merge_cols = [col for col in df.columns if col not in ["id", "updated_at"]]
agg_exprs = []

for col_name in merge_cols:
    # 收集字段的所有(时间,值)对,按时间倒序后过滤无效值,取第一个有效值
    valid_value = F.element_at(
        F.filter(
            F.array_sort(F.collect_list(F.struct(F.desc("updated_at").alias("ts"), col_name))),
            lambda x: x[col_name] != "<*no-update*>"
        ), 1
    )[col_name].alias(col_name)
    agg_exprs.append(valid_value)

# 补充取最新的更新时间
agg_exprs.append(F.max("updated_at").alias("updated_at"))

# 按ID分组聚合得到最终结果
result_df = df.groupBy("id").agg(*agg_exprs)
result_df.show()

输出结果

两种实现方式均会得到预期结果:

+---+-------+-------+----------+
| id|column1|column2|updated_at|
+---+-------+-------+----------+
|123|update2|     80|1634229000|
+---+-------+-------+----------+

内容的提问来源于stack exchange,提问作者Vipin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 10:45:02