如何用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
相关产品推荐
相关产品推荐

