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

如何在Apache Spark中实现CDC?新手求PySpark日数据差异提取方案

在Apache Spark(Python)中实现CDC并提取数据差异的方案

针对你作为Spark Python新手的需求,这里提供从基础批处理到进阶CDC的实用方案,先适配你当前的每日CSV批处理场景:

一、基础版:提取两日CSV数据的差异(新增/修改/删除)

假设你的数据有唯一主键(比如id),按以下步骤操作:

1. 读取并标记两日数据

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit

# 初始化SparkSession(若未创建)
spark = SparkSession.builder.appName("DailyDataDiff").getOrCreate()

# 读取前日数据(你已完成此步骤,这里补全示例)
prev_day_df = spark.read.csv("path/to/previous_day.csv", header=True, inferSchema=True)
# 读取当日数据
curr_day_df = spark.read.csv("path/to/current_day.csv", header=True, inferSchema=True)

# 为数据添加来源标识,方便后续区分
prev_day_df = prev_day_df.withColumn("source", lit("prev_day"))
curr_day_df = curr_day_df.withColumn("source", lit("curr_day"))

2. 全外连接并分类差异

通过全外连接匹配两日记录,筛选出三种核心差异类型:

# 按主键(如id)全外连接两日数据,添加后缀区分同名字段
joined_df = prev_day_df.join(curr_day_df, on="id", how="full_outer", suffixes=("_prev", "_curr"))

# 1. 新增记录:当日存在、前日不存在
new_records = joined_df.filter(col("id_prev").isNull()) \
                       .select([col(f"{c}_curr").alias(c) for c in curr_day_df.columns])

# 2. 删除记录:前日存在、当日不存在
deleted_records = joined_df.filter(col("id_curr").isNull()) \
                           .select([col(f"{c}_prev").alias(c) for c in prev_day_df.columns])

# 3. 修改记录:主键存在,但其他字段有差异
# 排除主键和来源列,生成字段比较条件
exclude_cols = ["id", "source"]
compare_cols = [c for c in prev_day_df.columns if c not in exclude_cols]
diff_conditions = [col(f"{c}_prev") != col(f"{c}_curr") for c in compare_cols]
full_diff_condition = diff_conditions[0]
for cond in diff_conditions[1:]:
    full_diff_condition = full_diff_condition | cond

modified_records = joined_df.filter(
    col("id_prev").isNotNull() & 
    col("id_curr").isNotNull() & 
    full_diff_condition
).select([col(f"{c}_curr").alias(c) for c in curr_day_df.columns]) \
 .withColumn("change_type", lit("modified"))

3. 输出差异结果

将三类差异数据保存为CSV或其他格式:

new_records.write.csv("path/to/new_records", header=True, mode="overwrite")
deleted_records.write.csv("path/to/deleted_records", header=True, mode="overwrite")
modified_records.write.csv("path/to/modified_records", header=True, mode="overwrite")

二、进阶版:长期CDC同步方案

若后续需要持续增量数据同步(而非每日手动处理CSV),可采用:

  • Spark CDC Connector:针对MySQL、PostgreSQL等数据库,直接捕获binlog生成增量数据,无需手动对比
  • Delta Lake:利用其版本控制特性,自动追踪数据变化,通过DESCRIBE HISTORY和deltaTable.asOfVersion()快速获取差异

关键注意事项

  • 必须确保数据有唯一主键,否则无法准确匹配两日记录
  • 读取CSV时尽量手动指定Schema(而非依赖inferSchema),避免类型不一致导致比较错误
  • 若字段数量较多,可封装辅助函数生成差异条件,减少重复代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 03:12:34