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

PySpark处理Delta表CDC数据:主键最新记录去重问题排查

问题描述

我正在开展一个Delta格式与数据表的探索项目,初始加载的是操作标识Op为'I'的记录,已通过PySpark读取并保存为Delta格式。后续处理变更数据捕获(CDC)文件时,在PySpark DataFrame中执行去重操作遇到问题,当前代码无法满足规则4的要求。

规则说明

  • 若某id仅含'I'操作记录,保留该记录;
  • 若某id最新记录为单个'U'操作,保留该记录;
  • 若某id存在多个'U'操作记录,保留sourceRecordTime最新的那条;
  • 若某id同时存在'I'和'U'操作记录,保留最新的'U'记录并将其Op改为'I'

当前代码、输出与期望输出

示例代码

from pyspark.sql.types import * 
from pyspark.sql.functions import to_timestamp, col, row_number

# 示例数据,元组列表
data = [
    ("I", "2024-03-22 22:49:56.000000", 71, 104.75),
    ("U", "2024-03-22 22:50:00.000000", 72, 114.75),
    ("I", "2024-03-22 22:49:56.000000", 73, 10.00),
    ("U", "2024-03-22 22:50:56.000000", 73, 20.00),
    ("U", "2024-03-22 22:51:56.000000", 73, 30.00),
    ("I", "2024-03-22 22:55:57.000000", 74, 30.00),
    ("U", "2024-03-22 22:49:56.000000", 75, 40.00),
    ("U", "2024-03-22 22:52:56.000000", 75, 50.00),
    ("U", "2024-03-22 22:57:56.000000", 75, 60.00)
]

# 定义Schema
schema = StructType([
    StructField("Op", StringType(), True),
    StructField("sourceRecordTime", StringType(), True),
    StructField("id", IntegerType(), True),
    StructField("amount", DoubleType(), True)
])

# 创建DataFrame
df = spark.createDataFrame(data, schema)
df = df.withColumn("sourceRecordTime", to_timestamp(df["sourceRecordTime"], "yyyy-MM-dd HH:mm:ss.SSSSSS"))

# 当前处理逻辑
sort_order = Window.partitionBy(col('id')).orderBy(col('sourceRecordTime').desc())
update_df = df.withColumn("rec_val", row_number().over(sort_order)).filter("rec_val=1").drop("rec_val")
update_df.show()

当前输出

+---+-------------------+---+------+
| Op|   sourceRecordTime| id|amount|
+---+-------------------+---+------+
|  I|2024-03-22 22:49:56| 71|104.75|
|  U|2024-03-22 22:50:00| 72|114.75|
|  U|2024-03-22 22:51:56| 73|  30.0|
|  I|2024-03-22 22:55:57| 74|  30.0|
|  U|2024-03-22 22:57:56| 75|  60.0|
+---+-------------------+---+------+

期望输出

+---+-------------------+---+------+
| Op|   sourceRecordTime| id|amount|
+---+-------------------+---+------+
|  I|2024-03-22 22:49:56| 71|104.75|
|  U|2024-03-22 22:50:00| 72|114.75|
|  I|2024-03-22 22:51:56| 73|  30.0|
|  I|2024-03-22 22:55:57| 74|  30.0|
|  U|2024-03-22 22:57:56| 75|  60.0|
+---+-------------------+---+------+
解决方案

当前代码仅筛选了每个id的最新记录,但未判断该id是否同时存在'I'和'U'操作,因此无法触发规则4的修改逻辑。我们需要先计算每个id的操作类型集合,再根据规则调整Op字段。

修改后的代码

from pyspark.sql.types import * 
from pyspark.sql.functions import to_timestamp, col, row_number, collect_set, when

# 示例数据与Schema定义(同原代码)
data = [
    ("I", "2024-03-22 22:49:56.000000", 71, 104.75),
    ("U", "2024-03-22 22:50:00.000000", 72, 114.75),
    ("I", "2024-03-22 22:49:56.000000", 73, 10.00),
    ("U", "2024-03-22 22:50:56.000000", 73, 20.00),
    ("U", "2024-03-22 22:51:56.000000", 73, 30.00),
    ("I", "2024-03-22 22:55:57.000000", 74, 30.00),
    ("U", "2024-03-22 22:49:56.000000", 75, 40.00),
    ("U", "2024-03-22 22:52:56.000000", 75, 50.00),
    ("U", "2024-03-22 22:57:56.000000", 75, 60.00)
]

schema = StructType([
    StructField("Op", StringType(), True),
    StructField("sourceRecordTime", StringType(), True),
    StructField("id", IntegerType(), True),
    StructField("amount", DoubleType(), True)
])

df = spark.createDataFrame(data, schema)
df = df.withColumn("sourceRecordTime", to_timestamp(df["sourceRecordTime"], "yyyy-MM-dd HH:mm:ss.SSSSSS"))

# 步骤1:计算每个id的操作类型集合,标记是否同时存在I和U
id_op_set_window = Window.partitionBy(col('id'))
df_with_op_set = df.withColumn(
    "op_types",
    collect_set(col('Op')).over(id_op_set_window)
)

# 步骤2:筛选每个id的最新记录,并根据规则调整Op字段
sort_order = Window.partitionBy(col('id')).orderBy(col('sourceRecordTime').desc())
final_df = df_with_op_set.withColumn("rec_val", row_number().over(sort_order)) \
    .filter("rec_val=1") \
    .drop("rec_val") \
    .withColumn(
        "Op",
        when(
            # 规则4:同时存在I和U,且当前记录是U时,改为I
            (col('op_types') == collect_set(['I', 'U'])) & (col('Op') == 'U'),
            'I'
        ).otherwise(col('Op'))
    ) \
    .drop("op_types")

final_df.show()

代码解释

  1. 计算操作类型集合:通过collect_set窗口函数,为每个id收集所有出现过的Op类型,以此判断是否同时存在'I'和'U'。
  2. 筛选最新记录:保留每个id的最新记录(同原逻辑)。
  3. 调整Op字段:使用when条件函数,当某个id同时存在'I'和'U'且最新记录是'U'时,将Op改为'I',其他情况保持原Op不变。

验证输出

运行修改后的代码,输出将与期望输出一致:

+---+-------------------+---+------+
| Op|   sourceRecordTime| id|amount|
+---+-------------------+---+------+
|  I|2024-03-22 22:49:56| 71|104.75|
|  U|2024-03-22 22:50:00| 72|114.75|
|  I|2024-03-22 22:51:56| 73|  30.0|
|  I|2024-03-22 22:55:57| 74|  30.0|
|  U|2024-03-22 22:57:56| 75|  60.0|
+---+-------------------+---+------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 23:37:02