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

如何使用Spark实现数据库CDC并将数据以Parquet格式写入HDFS

基于Spark实现MySQL CDC并写入HDFS Parquet的解决方案

你当前的代码属于全量同步逻辑,要实现变更数据捕获,可根据业务场景选择以下两种实现方案:


方案1:基于增量字段的轻量CDC实现

适用于小数据量、源表存在可识别增量特征字段的场景,无需额外部署组件。

前提条件

源表需要具备更新时间字段(如update_time,行新增/修改时自动更新)或者自增主键,用于筛选每次同步的变更数据。

改造后核心逻辑

from pyspark.sql import SparkSession
import os

spark = SparkSession \
        .builder \
        .appName("Incremental_Ingest") \
        .master("local[*]") \
        .config("spark.driver.extraClassPath", "/home.../mysql-connector-java-5.1.30.jar") \
        .getOrCreate()

# 1. 读取上次同步的最大位点(可存在HDFS的标记文件/配置中心,首次同步设为'1970-01-01 00:00:00')
last_sync_time = "2024-01-01 00:00:00"
if os.path.exists("/tmp/last_sync_time.txt"):
    with open("/tmp/last_sync_time.txt", "r") as f:
        last_sync_time = f.read().strip()

# 2. 增量拉取数据,加where条件避免全表扫描
df = spark.read\
        .format("jdbc") \
        .option("url", "jdbc:mysql://localhost:3306/classicmodels") \
        .option("driver", "com.mysql.jdbc.Driver") \
        .option("dbtable", f"(select * from employees where update_time >= '{last_sync_time}') as tmp") \
        .option("user", "...") \
        .option("password", "...").load()

# 3. 写入HDFS Parquet,开启snappy压缩节省空间
df.write.mode("append")\
        .option("compression", "snappy")\
        .parquet("hdfs://localhost:9000/employees_parquet/")

# 4. 更新本次同步的最大位点,供下次同步使用
current_max_time = df.agg({"update_time": "max"}).collect()[0][0]
with open("/tmp/last_sync_time.txt", "w") as f:
    f.write(str(current_max_time))

注意事项

  • 如果需要捕获删除事件,该方案无法支持,需要配合逻辑删除字段(如is_deleted)使用
  • 写入时如果有更新合并需求,可以搭配Iceberg/Delta Lake等湖存储格式实现Upsert,避免数据重复

方案2:基于Binlog的实时CDC实现

适用于大数据量、低延迟、需要捕获全量变更(新增/修改/删除)的场景。

核心流程

  1. 开启MySQL Binlog,格式设置为ROW模式
  2. 部署Debezium连接器监听目标表,将所有变更事件写入Kafka
  3. 用Spark Structured Streaming消费Kafka中的CDC事件,解析后写入HDFS

示例代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, date_format
from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType

spark = SparkSession \
    .builder \
    .appName("Realtime_CDC_Ingest") \
    .getOrCreate()

# 定义Debezium返回的CDC事件结构,需和实际表字段匹配
cdc_schema = StructType([
    StructField("op", StringType(), True), # 变更类型:c=新增,u=更新,d=删除
    StructField("before", StructType([
        StructField("id", LongType(), True),
        StructField("name", StringType(), True),
        StructField("update_time", TimestampType(), True)
    ]), True),
    StructField("after", StructType([
        StructField("id", LongType(), True),
        StructField("name", StringType(), True),
        StructField("update_time", TimestampType(), True)
    ]), True),
    StructField("ts_ms", LongType(), True)
])

# 消费Kafka中的CDC数据
raw_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "mysql_classicmodels_employees_cdc") \
    .load()

# 解析CDC数据并按天分区
parsed_df = raw_df.select(from_json(col("value").cast("string"), cdc_schema).alias("cdc")) \
    .select(
        "cdc.op",
        "cdc.before",
        "cdc.after",
        col("cdc.ts_ms").cast(TimestampType()).alias("event_time")
    ).withColumn("dt", date_format("event_time", "yyyy-MM-dd"))

# 流式写入HDFS Parquet,配置checkpoint保证 Exactly-Once 语义
query = parsed_df.writeStream \
    .format("parquet") \
    .option("path", "hdfs://localhost:9000/realtime_cdc/employees/") \
    .option("checkpointLocation", "hdfs://localhost:9000/cdc_checkpoint/employees/") \
    .option("compression", "snappy") \
    .partitionBy("dt") \
    .trigger(processingTime="5 minutes") \
    .start()

query.awaitTermination()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 03:27:01