Flink CDC 实时同步:MySQL 到 MongoDB 实战指南

以下步骤基于 Flink 1.14+Flink CDC 2.3+,5 分钟内完成从 MySQL 到 MongoDB 的实时同步。


1. 环境准备(1分钟)
  • 依赖 Jar 包:下载并放入 Flink 的 lib/ 目录:
    • flink-sql-connector-mysql-cdc-2.3.0.jar
    • flink-connector-mongodb-1.0.0.jar
  • 数据库配置
    • MySQL:开启 binlog(在 my.cnf 中添加):
      server-id=1
      log_bin=mysql-bin
      binlog_format=ROW
      

    • MongoDB:无需特殊配置。

2. 编写 Flink SQL 同步脚本(2分钟)

在 Flink SQL Client 中执行以下代码(替换 <> 内参数):

-- 创建 MySQL CDC 源表
CREATE TABLE mysql_users (
    id INT PRIMARY KEY,
    name STRING,
    email STRING
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '<mysql_host>',
    'port' = '3306',
    'username' = '<user>',
    'password' = '<password>',
    'database-name' = '<db_name>',
    'table-name' = 'users'
);

-- 创建 MongoDB Sink 表
CREATE TABLE mongodb_users (
    id INT PRIMARY KEY,
    name STRING,
    email STRING
) WITH (
    'connector' = 'mongodb',
    'uri' = 'mongodb://<mongo_host>:27017',
    'database' = '<mongo_db>',
    'collection' = 'users'
);

-- 启动实时同步
INSERT INTO mongodb_users 
SELECT * FROM mysql_users;


3. 验证同步(1分钟)
  1. 在 MySQL 插入数据
    INSERT INTO users VALUES (1, 'Alice', 'alice@example.com');
    

  2. 检查 MongoDB
    use <mongo_db>
    db.users.find({id: 1})  -- 应返回 Alice 的记录
    


4. 关键配置说明
  • MySQL CDC 参数
    • 'scan.startup.mode' = 'latest-offset':从最新变更开始同步(支持 initial 全量+增量)。
  • MongoDB Sink 参数
    • 'sink.buffer-flush.max-rows' = '100':每批次写入最大行数。
    • 'sink.buffer-flush.interval' = '10s':写入间隔。

5. 常见问题
  • 数据延迟:检查 Flink 任务是否积压(Web UI 的 BackPressure 标签)。
  • 同步失败:确认 MySQL 的 binlog 权限和 MongoDB 网络连通性。
  • 字段映射:源表和目标表字段需名称一致(或通过 AS 重命名)。

✅ 完成!从 MySQL 插入到 MongoDB 写入延迟通常在 1-3 秒内。

更多推荐