Flink CDC 3.0 实战:零干预实现MySQL到Doris的整库实时同步

最近在重构一个老旧的报表系统,数据源来自几个核心的MySQL业务库,表结构隔三差五就要调整一版。每次业务方加个字段,我们这边就得停掉同步任务,手动改下游Doris表结构,再重启任务,运维同学都快被搞疯了。直到我们开始尝试Flink CDC 3.0,情况才彻底改变。这个版本带来的整库同步动态表结构变更处理能力,简直是为这类场景量身定做的。今天,我就结合实际的踩坑经验,带你从零开始,快速搭建一个能自动应对源表结构变化的实时数据同步管道。

1. 环境准备与核心组件部署

在开始编写任何同步任务之前,一个稳定、版本匹配的运行环境是基础。Flink CDC 3.0 对上下游组件的版本有一定要求,配置不当很容易在后期遇到兼容性问题。

1.1 Flink 集群部署与关键配置

我们选择 Flink 1.18.0 作为运行引擎,这是目前与 Flink CDC 3.0 兼容性最好的稳定版本。下载解压后,首要任务是调整 conf/flink-conf.yaml 文件。对于CDC任务,以下几个参数至关重要:

# 启用并配置Checkpoint,这是保证Exactly-Once语义和任务故障恢复的基石
execution.checkpointing.interval: 3000ms
execution.checkpointing.mode: EXACTLY_ONCE
execution.checkpointing.timeout: 10min

# 根据你的服务器资源调整TaskManager的Slot数量
# 一个典型的整库同步任务可能需要多个Slot来处理并行读取和写入
taskmanager.numberOfTaskSlots: 4

# 状态后端建议使用RocksDB,以应对可能较大的状态数据
state.backend: rocksdb
state.checkpoints.dir: file:///tmp/flink-checkpoints

配置完成后,使用 ./bin/start-cluster.sh 启动集群。通过 http://localhost:8081 访问Web UI,确认所有组件(JobManager、至少一个TaskManager)运行正常。一个常见的误区是只启动一个TaskManager且Slot数不足,提交任务时会报资源不足的错误。你可以通过多次执行 start-cluster.sh 或在配置文件中增加 taskmanager.numberOfTaskSlots 来扩容。

1.2 使用Docker Compose快速搭建MySQL与Doris

为了演示的纯粹性,我们使用Docker Compose在本地一键拉起MySQL和Doris服务。这能避免环境差异带来的问题。

注意:Doris的BE节点依赖于Linux系统的内存映射配置。在宿主机上执行 sudo sysctl -w vm.max_map_count=2000000 是必要的,否则Doris容器可能无法启动。

创建一个 docker-compose.yml 文件,内容如下:

version: '3.8'
services:
  mysql-cdc-demo:
    image: mysql:8.0
    container_name: mysql-cdc-demo
    environment:
      MYSQL_ROOT_PASSWORD: 123456
      MYSQL_DATABASE: app_db
    ports:
      - "3306:3306"
    command: 
      --server-id=1 
      --log-bin=mysql-bin 
      --binlog-format=ROW 
      --gtid-mode=ON 
      --enforce-gtid-consistency=ON

  doris-fe:
    image: apache/doris:1.2.4-fe-x86_64
    container_name: doris-fe
    environment:
      FE_SERVERS: "fe1:172.20.80.1:9010"
      FE_ID: 1
    ports:
      - "8030:8030"  # HTTP端口,用于Web UI和连接
      - "9030:9030"  # MySQL协议端口,用于查询
    volumes:
      - ./doris-data/fe:/opt/doris/fe/doris-meta

  doris-be:
    image: apache/doris:1.2.4-be-x86_64
    container_name: doris-be
    environment:
      FE_SERVERS: "fe1:172.20.80.1:9010"
      BE_ADDR: "172.20.80.2:9050"
    depends_on:
      - doris-fe
    volumes:
      - ./doris-data/be:/opt/doris/be/storage
    sysctls:
      - vm.max_map_count=2000000

这个配置做了几件关键事:

  1. MySQL:启用了基于GTID的二进制日志(binlog),这是CDC捕获变更数据的前提。
  2. Doris:分别启动了FE(前端)和BE(后端)服务,并配置了服务发现。数据卷挂载是为了持久化数据。

在文件所在目录执行 docker-compose up -d,然后通过 docker-compose logs -f 观察日志,直到所有服务显示健康状态。你可以通过 mysql -h 127.0.0.1 -P 3306 -uroot -p123456 连接MySQL,并通过浏览器访问 http://localhost:8030(用户名root,密码为空)进入Doris Web UI。

2. 理解Flink CDC 3.0的核心机制

在动手写配置之前,花点时间理解Flink CDC 3.0的工作原理,能让你在遇到问题时更快地定位和解决。3.0版本最大的革新在于引入了 Schema Registry(模式注册中心)Schema Operator(模式操作器) 这两个核心组件。

想象一下传统CDC同步的流程:Source端读取数据,经过一些转换,然后Sink端写入。当源表增加一个字段时,这条新增字段的数据会流入管道,但Sink端的表结构还没变,于是写入失败,任务崩溃。

Flink CDC 3.0的解决方案非常巧妙。它在作业拓扑中增加了一个“交通警察”——Schema Registry。当MySQL端发生 ALTER TABLE 操作时,这个DDL事件会作为一个特殊的“模式变更事件”被捕获,并首先发送到Schema Registry。Schema Operator会介入,它的工作流程可以概括为:

  1. 暂停数据流:暂时停止普通数据记录向下游流动。
  2. 清空流水线:确保所有在途的、基于旧模式的数据都被Sink端完全处理并写出。
  3. 协调模式变更:将DDL事件发送到下游Sink(如Doris),驱动其完成表结构变更。
  4. 恢复数据流:当下游确认结构变更成功后,恢复数据流动,此后所有新数据(包括新字段)都将按新结构处理。

这个过程保证了模式变更的原子性和一致性,实现了真正的“动态表结构变更同步”。对于整库同步,其本质是利用正则表达式匹配库表名,为匹配到的每一张表动态创建和维护一个独立的CDC读取链路和写入链路,并通过统一的Job进行管理。

3. 从零编写整库同步任务

理论清晰后,我们进入实战环节。Flink CDC 3.0提供了两种任务提交方式:通过其自带的 flink-cdc.sh CLI工具,或者将CDC Connector作为库引入到自定义的Flink SQL/DataStream作业中。对于快速启动和运维,CLI方式更加简洁。

3.1 获取与配置Connector

首先,从官网下载 flink-cdc-3.0.0-bin.tar.gz 并解压。然后,下载以下两个必须的Connector JAR包,放入解压后目录的 lib/ 文件夹下:

  • flink-sql-connector-mysql-cdc-3.0.0.jar
  • doris-connector-1.16-3.0.0.jar (注意选择与你目标Flink版本匹配的Doris Connector)

提示:Connector版本必须严格匹配。使用错误的版本可能会导致无法识别的配置参数或序列化错误。

3.2 编写核心配置文件

接下来,创建任务配置文件 mysql-to-doris.yaml。这个YAML文件定义了从源到目标的完整数据管道。

################################################################################
# 描述:将MySQL整库同步至Doris,并自动处理表结构变更
################################################################################
source:
  type: mysql
  hostname: localhost
  port: 3306
  username: root
  password: 123456
  database: app_db
  table-list: [“app_db.*”] # 使用通配符匹配库下所有表
  server-id: 5400-5404     # 为每个表读取任务分配独立的server-id范围,避免冲突
  server-time-zone: Asia/Shanghai
  scan.startup.mode: initial # 首次启动时先做全量快照,再接增量流

sink:
  type: doris
  fenodes: localhost:8030
  username: root
  password: ""
  table.create.properties.light_schema_change: true # 启用轻量表结构变更
  table.create.properties.replication_num: 1        # 单机Doris演示,副本数设为1

pipeline:
  name: MySQL_to_Doris_Whole_Database_Sync
  parallelism: 2 # 根据表数量和复杂度调整并行度

关键配置解析:

  • source.table-list: 这里的 app_db.* 是整库同步的“魔法咒语”。它会自动发现 app_db 库下的所有现有表,并为每张表启动一个独立的CDC读取子任务。
  • source.server-id: 在MySQL主从复制中,每个slave需要一个唯一的ID。Flink CDC作为MySQL的“逻辑从库”,也需要这个ID。当并行读取多张表时,需要提供一个ID范围,框架会自动分配。
  • sink.table.create.properties.light_schema_change: 这是Doris Connector的关键参数。设置为 true 后,Connector在接收到新增列的DDL事件时,会尝试使用Doris的 Light Schema Change 功能来添加列,这通常是一个毫秒级的轻量操作。如果设为 false 或不设置,对于某些变更可能需要重写整个数据文件,代价巨大。
  • scan.startup.mode: initial 模式是默认且最常用的,它先读取表的当前全量数据快照,然后无缝切换到读取binlog进行增量变更。这对于初始化同步场景非常友好。

3.3 提交任务与验证

flink-cdc-3.0.0 目录下,执行以下命令提交任务:

./bin/flink-cdc.sh conf/mysql-to-doris.yaml

如果一切顺利,命令行会输出提交成功的Job ID。立刻打开Flink Web UI (localhost:8081),你应该能看到一个名为 MySQL_to_Doris_Whole_Database_Sync 的任务在运行。

现在,回到MySQL客户端,创建一些测试表和初始数据:

USE app_db;
CREATE TABLE user_behavior (
    user_id BIGINT,
    item_id BIGINT,
    category_id INT,
    behavior VARCHAR(20),
    ts TIMESTAMP,
    PRIMARY KEY(user_id, item_id, ts)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

INSERT INTO user_behavior VALUES 
(1001, 2001, 1, 'click', '2024-01-01 10:00:00'),
(1001, 2002, 2, 'buy', '2024-01-01 10:01:00');

CREATE TABLE product_info (
    product_id INT PRIMARY KEY,
    product_name VARCHAR(255),
    price DECIMAL(10, 2),
    update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);

INSERT INTO product_info (product_id, product_name, price) VALUES 
(1, 'Laptop', 5999.99),
(2, 'Mouse', 89.50);

稍等片刻(取决于数据量),刷新Doris Web UI的 app_db 数据库,你会发现 user_behaviorproduct_info 两张表已经被自动创建,并且初始数据已经同步完成。这个过程完全无需你手动在Doris中执行任何 CREATE TABLE 语句。

4. 应对动态表结构变更与复杂场景

静态同步只是开始,真正的考验在于变化。我们模拟几个业务中常见的场景。

4.1 体验自动化Schema Evolution

现在,业务方要求在 user_behavior 表中增加一个 device 字段来记录用户使用的设备类型。在旧模式下,你需要:

  1. 暂停Flink作业。
  2. 在Doris中执行 ALTER TABLE app_db.user_behavior ADD COLUMN device VARCHAR(32) NULL;
  3. 重启Flink作业。

而现在,你只需要在MySQL端执行:

ALTER TABLE app_db.user_behavior ADD COLUMN device VARCHAR(32) NULL COMMENT ‘用户设备’;

执行完成后,立即向表中插入一条包含新字段的数据:

INSERT INTO user_behavior VALUES (1002, 2003, 3, ‘pv’, ‘2024-01-01 11:00:00’, ‘iOS’);

此时,观察Flink Web UI中任务的管理日志(或通过 docker-compose logs -f 观察Doris FE日志),你会看到类似“Schema change event received”、“Applying ALTER to Doris”的信息流。很快,再次查询Doris中的 user_behavior 表,你会发现:

  1. 表结构已经自动增加了 device 列。
  2. 新插入的那条完整数据(包含device=’iOS’)已经存在。
  3. 任务从未停止,数据同步的延迟仅在Schema变更的瞬间有轻微增加。

这就是 “零手工干预” 的魅力。无论是增加列、删除列(需注意Doris对删除列的支持度)、修改列类型(需兼容),还是重命名列,整个过程都可以自动化完成。

4.2 实现分库分表数据汇聚

另一个经典场景是分库分表。例如,订单数据按月份分表,如 order_202401, order_202402。业务上希望将这些表同步到Doris后,合并到一张名为 ods_order_all 的宽表中进行分析。

Flink CDC 3.0 的 路由(Route) 功能正是为此而生。你需要修改配置文件,添加 route 配置节:

source:
  ...
  table-list: [“app_db.order_202401”, “app_db.order_202402”, “app_db.order_202403”] # 明确列出需要同步的分表
  ...

sink:
  ...
  # 注意,这里不再需要指定表名,路由规则会覆盖

route:
  - source-table: “app_db\.order_.*”   # 使用正则匹配所有月份订单表
    sink-table: “app_db.ods_order_all” # 统一路由到目标表

这里有一个至关重要的细节:正则表达式转义。 在YAML中,反斜杠 \ 本身是转义字符。为了表示正则表达式中的 \.(匹配点号),你必须写成 \\.。上面配置中的 app_db\\.order_.* 会被正确解析为正则 app_db\.order_.*,从而匹配 app_db.order_202401 等表。如果写成 app_db.order_.*,点号在正则中代表“任意单个字符”,会导致匹配过度或错误。

提交这个带路由配置的任务后,所有源分表的数据都会汇聚写入到Doris的 ods_order_all 这一张表中。这对于下游的BI分析来说,透明地解决了数据分散的问题。

4.3 路由场景下的Schema变更挑战与策略

然而,分库分表同步遇到动态Schema变更时,情况会变得复杂。假设 order_202401 表增加了一个 coupon_amount 字段,而 order_202402 表没有变。根据我们之前了解的机制,Schema Operator会尝试将 ADD COLUMN coupon_amount 的DDL事件应用到目标表 ods_order_all 上。这本身是成功的。

问题出在后续的数据流上。当 order_202402 表的数据(不包含新字段)到达Sink端时,其字段数与目标表当前的结构不匹配,就会抛出 Column size does not match the data size 这类异常。

注意:这是目前Flink CDC 3.0在分库分表合并同步场景下的一个限制。官方文档也指出,当前版本在多源表结构不一致时,向同一个目标表同步可能会遇到问题。

应对策略:

  1. 业务规范先行:确保进行合并同步的多个源分表,其Schema变更(如加字段)必须同步进行。这通常需要靠业务开发规范和上线流程来保证。
  2. 使用物化视图或后续加工:如果无法保证同时变更,可以考虑先将各分表同步到Doris的不同独立表中(不使用路由合并),然后在Doris内部通过物化视图(Materialized View) 或定时调度任务进行合并。这样,每个独立表的Schema可以独立演进,物化视图在查询时再处理字段对齐问题。
  3. 期待社区版本更新:Apache Flink社区已经意识到这个痛点,并在后续版本规划中致力于提供更优雅的多源Schema合并解决方案。

5. 生产级部署考量与监控

将演示项目推向生产环境,还需要考虑更多因素。

高可用与容错:文中的Standalone集群模式不适合生产。应部署在 Flink on YARN/K8s 上,并配置高可用的Checkpoint存储(如HDFS、S3)。确保 state.backendhigh-availability 相关配置正确。

性能调优:整库同步大量表时,对MySQL源端的连接数和binlog读取压力需要评估。

  • server-id 范围要足够大。
  • 适当调整 parallelism,但注意不要超过Flink集群的总Slot数。
  • 对于数据量巨大的单表,可以配置 scan.incremental.snapshot.chunk.size 来调整全量快照阶段的分片大小,平衡内存使用和速度。

监控与告警:除了Flink Web UI,应集成到现有的监控体系(如Prometheus + Grafana)。需要重点关注以下指标:

  • 源端延迟currentFetchEventTimeLag,表示处理的数据与当前时间的时间差,是衡量同步实时性的核心指标。
  • Checkpoint健康状况:成功率、耗时、大小。频繁的Checkpoint失败通常意味着背压或状态过大。
  • Doris写入指标:通过Doris的监控接口,关注导入频率、失败次数和耗时。

一个实用的经验是,在首次启动全库同步任务前,如果源库已有大量历史数据,最好先在业务低峰期启动,并密切关注Flink TaskManager的内存使用情况以及MySQL的负载。你可以考虑先通过 scan.startup.mode: latest-offset 模式启动,只接增量流,然后另用离线工具初始化历史数据,最后再切换任务模式进行校验和追增,以降低对线上数据库的冲击。

从我的实践来看,Flink CDC 3.0的整库同步功能已经极大地简化了数据入仓的初始化和日常运维工作。尤其是在应对频繁的敏捷开发迭代时,数据团队不再需要被动地响应每一个表结构变更请求。当然,像分表合并同步这类复杂场景,还需要结合具体的业务约束和技术变通方案。建议在全面铺开前,针对核心业务场景进行充分的测试,特别是模拟各种异常Schema变更和网络抖动的情况,确保整个管道的健壮性符合你的SLA要求。

更多推荐