如何使用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实现
适用于大数据量、低延迟、需要捕获全量变更(新增/修改/删除)的场景。
核心流程
- 开启MySQL Binlog,格式设置为ROW模式
- 部署Debezium连接器监听目标表,将所有变更事件写入Kafka
- 用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
相关产品推荐
相关产品推荐

