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

