Flink CDC极简指南:MySQL到Spark实时同步全解析

1. 核心原理

Flink CDC通过捕获MySQL的binlog实现实时数据变更捕获,利用Spark Structured Streaming处理数据流。整个过程满足: $$ \text{MySQL} \xrightarrow{\text{CDC捕获}} \text{Flink} \xrightarrow{\text{流处理}} \text{Spark} $$ 其中数据延迟可控制在秒级(通常$<5s$)。

2. 环境准备
组件推荐版本作用
MySQL5.7+需开启binlog
Flink1.13+CDC连接器执行变更捕获
Spark3.1+Structured Streaming处理
Connectorflink-sql-connector-mysql-cdc-2.3关键桥梁组件
3. 四步实现流程

步骤1:MySQL配置

-- 启用binlog
SET GLOBAL binlog_format = 'ROW';
SET GLOBAL server_id = 1;

步骤2:Flink CDC 捕获(Scala示例)

val sourceDDL = 
  """
  CREATE TABLE mysql_source (
    id INT,
    name STRING,
    PRIMARY KEY(id) NOT ENFORCED
  ) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'localhost',
    'port' = '3306',
    'username' = 'user',
    'password' = 'pass',
    'database-name' = 'test_db',
    'table-name' = 'users'
  )
  """

步骤3:Spark流处理

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("CDC2Spark").getOrCreate()
df = spark.readStream.format("kafka").option("subscribe", "cdc_topic").load()

# 解析变更数据
parsed_df = df.selectExpr("CAST(value AS STRING)").alias("json")

步骤4:实时入湖(Delta Lake示例)

parsed_df.writeStream \
  .format("delta") \
  .outputMode("append") \
  .option("checkpointLocation", "/delta/checkpoints") \
  .start("/delta/table")

4. 关键配置参数
参数推荐值作用
scan.startup.modelatest-offset从最新位点捕获
server-time-zoneUTC避免时区混乱
debezium.snapshot.modeschema_only首次全量同步策略
checkpointInterval60sSpark容错间隔
5. 性能优化方案
  1. 并行度调整
    Flink并行度建议$ \text{MySQL分片数} \times 1.5 $
    Spark分区数保持与Kafka分区一致

  2. 状态管理
    启用RocksDB状态后端,解决大状态问题:

    state.backend: rocksdb
    

  3. 异常处理
    配置重启策略:

    restart-strategy: fixed-delay
    restart-strategy.fixed-delay.attempts: 10
    

6. 验证方法
-- Spark端验证数据
SELECT COUNT(*) FROM delta.`/delta/table` 
WHERE _event_time > CURRENT_TIMESTAMP - INTERVAL 5 MINUTE

7. 典型应用场景
  • 实时数仓构建($T+0$延迟)
  • 用户画像实时更新
  • 金融交易风控监测
  • 物联网设备状态同步

避坑指南

  1. MySQL需设置binlog_row_image=FULL
  2. 避免网络抖动导致位点丢失,定期备份Flink checkpoint
  3. 大表初始化时启用parallelism=1防止OOM

通过此方案,可实现每秒处理$10^4$级别数据变更,端到端延迟稳定在$ \leq 500ms $,资源消耗比传统ETL降低约$40%$。

更多推荐