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()
代码解释
- 计算操作类型集合:通过
collect_set窗口函数,为每个id收集所有出现过的Op类型,以此判断是否同时存在'I'和'U'。 - 筛选最新记录:保留每个
id的最新记录(同原逻辑)。 - 调整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
相关产品推荐
相关产品推荐

