返回博客

Flink CDC 实时数据同步指南

在企业数据平台中,数据同步是最基础也是最关键的能力之一。传统的批量 ETL 虽然稳定,但无法满足实时决策的需求。Flink CDC(Change Data Capture)让你可以实时捕获数据库的变更事件,将其流式传输到 coomia-dip 的数据层,让业务数据库中的每一行变更都在毫秒级别反映到分析平台中。

Coomia发布于 2026年1月25日8 分钟阅读
分享本文Twitter / X

系列:S12 开发者教程 · 第 15 篇 | 难度:中级 | 阅读时间:15 分钟

Flink CDC 实时数据同步指南

#引言

在企业数据平台中,数据同步是最基础也是最关键的能力之一。传统的批量 ETL 虽然稳定,但无法满足实时决策的需求。Flink CDC(Change Data Capture)让你可以实时捕获数据库的变更事件,将其流式传输到 coomia-dip 的数据层,让业务数据库中的每一行变更都在毫秒级别反映到分析平台中。

本教程将带你理解 CDC 的原理,配置 MySQL binlog,部署 Flink 集群,并实现一个从 MySQL 到 Apache Doris 的实时同步管道。

#1. CDC 基础概念

#1.1 什么是 Change Data Capture

Change Data Capture 是一种通过监听数据库事务日志来捕获数据变更的技术。与定时全量拉取相比,CDC 有三个核心优势:

  • 实时性:变更发生后毫秒级捕获,而非等待下一个调度周期
  • 低负载:只传输变更数据,不需要全表扫描
  • 完整性:捕获所有变更(包括中间状态),而非只有最终快照

CDC 的工作流程可以简化为:源数据库的事务日志被 Flink CDC Connector 持续消费,解析为结构化的变更事件(Changelog Stream),经过 Flink 的流处理进行清洗、转换和路由,最终写入目标系统(在 coomia-dip 中是 Apache Doris)。

Flink CDC 3.x 引入了声明式 Pipeline 架构,大幅简化了同步任务的配置。你只需要一个 YAML 文件即可描述整个同步拓扑:

YAML
source:
  type: mysql
  hostname: mysql-source
  port: 3306
  username: cdc_user
  password: ${MYSQL_CDC_PASSWORD}
  tables: erp_db.orders, erp_db.order_items, erp_db.customers

sink:
  type: doris
  fenodes: doris-fe:8030
  username: root
  password: ${DORIS_PASSWORD}

pipeline:
  name: erp-to-doris-sync
  parallelism: 4
  schema.change.behavior: evolve

#1.3 coomia-dip 中的 CDC 定位

在 coomia-dip 的八平面架构中,CDC 属于 Data Layer(Data Layer)的核心能力。数据从业务数据库通过 CDC 进入 Doris ODS 层,然后由 Ontology Mapper 自动映射为 Ontology 对象,最终通过 OQL 引擎提供统一查询。

#2. 环境准备

#2.1 启用 MySQL Binlog

CDC 依赖 MySQL 的 binlog(二进制日志)。必须使用 ROW 格式以获取完整的字段级变更信息:

INI
[mysqld]
server-id         = 1
log_bin           = mysql-bin
binlog_format     = ROW
binlog_row_image  = FULL
expire_logs_days  = 3
gtid_mode         = ON
enforce_gtid_consistency = ON

创建 CDC 专用数据库用户:

SQL
CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'cdc_secure_password';
GRANT SELECT, RELOAD, SHOW DATABASES,
      REPLICATION SLAVE, REPLICATION CLIENT
    ON *.* TO 'cdc_user'@'%';
FLUSH PRIVILEGES;

验证配置:

SQL
SHOW VARIABLES LIKE 'log_bin';          -- ON
SHOW VARIABLES LIKE 'binlog_format';    -- ROW
SHOW VARIABLES LIKE 'binlog_row_image'; -- FULL
SHOW MASTER STATUS;

使用 Docker Compose 快速部署 Flink Session 集群:

YAML
version: "3.8"
services:
  jobmanager:
    image: flink:1.18-java11
    ports: ["8081:8081"]
    command: jobmanager
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
        state.checkpoints.dir: file:///opt/flink/checkpoints
        execution.checkpointing.interval: 60000
        execution.checkpointing.mode: EXACTLY_ONCE

  taskmanager:
    image: flink:1.18-java11
    depends_on: [jobmanager]
    command: taskmanager
    deploy:
      replicas: 2
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
        taskmanager.numberOfTaskSlots: 4
        taskmanager.memory.process.size: 4096m

#2.3 安装 CDC Connector

Bash
FLINK_LIB=/opt/flink/lib

wget -P $FLINK_LIB \
  https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/3.0.1/flink-sql-connector-mysql-cdc-3.0.1.jar

wget -P $FLINK_LIB \
  https://repo1.maven.org/maven2/org/apache/doris/flink-doris-connector-1.18/1.6.1/flink-doris-connector-1.18-1.6.1.jar

#3. 实现同步管道

Flink SQL 是最快捷的 CDC 上手方式,不需要编写 Java 代码:

SQL
SET 'execution.runtime-mode' = 'streaming';
SET 'execution.checkpointing.interval' = '60s';

CREATE TABLE orders_source (
    order_id     BIGINT,
    customer_id  BIGINT,
    product_name STRING,
    quantity     INT,
    unit_price   DECIMAL(10, 2),
    total_amount DECIMAL(12, 2),
    status       STRING,
    created_at   TIMESTAMP(3),
    updated_at   TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector'                             = 'mysql-cdc',
    'hostname'                              = 'mysql-source',
    'port'                                  = '3306',
    'username'                              = 'cdc_user',
    'password'                              = 'cdc_secure_password',
    'database-name'                         = 'erp_db',
    'table-name'                            = 'orders',
    'scan.incremental.snapshot.enabled'      = 'true',
    'scan.incremental.snapshot.chunk.size'   = '8096'
);

CREATE TABLE orders_sink (
    order_id     BIGINT,
    customer_id  BIGINT,
    product_name STRING,
    quantity     INT,
    unit_price   DECIMAL(10, 2),
    total_amount DECIMAL(12, 2),
    status       STRING,
    created_at   TIMESTAMP(3),
    updated_at   TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector'                           = 'doris',
    'fenodes'                             = 'doris-fe:8030',
    'table.identifier'                    = 'ods.orders',
    'username'                            = 'root',
    'password'                            = '',
    'sink.properties.format'              = 'json',
    'sink.properties.read_json_by_line'   = 'true',
    'sink.enable-delete'                  = 'true',
    'sink.label-prefix'                   = 'flink_cdc_orders'
);

INSERT INTO orders_sink SELECT * FROM orders_source;

#3.2 DataStream API(灵活定制)

对于需要复杂转换逻辑的场景:

Java
package com.onto.pipeline.cdc;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;

public class ERPSyncPipeline {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
        env.setParallelism(4);

        MySqlSource<String> source = MySqlSource.<String>builder()
            .hostname("mysql-source")
            .port(3306)
            .databaseList("erp_db")
            .tableList("erp_db.orders", "erp_db.order_items", "erp_db.customers")
            .username("cdc_user")
            .password("cdc_secure_password")
            .deserializer(new JsonDebeziumDeserializationSchema())
            .includeSchemaChanges(true)
            .splitSize(8096)
            .build();

        env.fromSource(source, WatermarkStrategy.noWatermarks(), "MySQL CDC")
            .map(new CDCEventTransformer())
            .process(new TableRoutingFunction())
            .sinkTo(DorisSinkBuilder.build());

        env.execute("ERP to coomia-dip CDC Pipeline");
    }
}

转换函数示例:

Java
public class CDCEventTransformer implements MapFunction<String, CDCRecord> {
    @Override
    public CDCRecord map(String json) throws Exception {
        JSONObject event = JSON.parseObject(json);
        String op = event.getString("op");
        JSONObject after = event.getJSONObject("after");
        JSONObject before = event.getJSONObject("before");
        JSONObject src = event.getJSONObject("source");

        CDCRecord record = new CDCRecord();
        record.setDatabase(src.getString("db"));
        record.setTable(src.getString("table"));
        record.setOperation(op);
        record.setData(after != null ? after : before);
        record.setTimestamp(System.currentTimeMillis());

        if ("orders".equals(record.getTable()) && after != null) {
            double amount = after.getDoubleValue("total_amount");
            after.put("amount_with_tax", amount * 1.13);
            after.put("currency", "CNY");
        }
        return record;
    }
}

#4. 与 coomia-dip Ontology 集成

#4.1 注册 CDC 数据源

Python
from ontology_sdk import OntoPlatform

platform = OntoPlatform(base_url="http://localhost:8080", token="admin-token")

source = platform.data_sources.register(
    name="erp-mysql-cdc",
    source_type="FLINK_CDC",
    config={
        "connector": "mysql-cdc",
        "hostname": "mysql-source",
        "port": 3306,
        "database": "erp_db",
        "tables": ["orders", "order_items", "customers"],
        "sync_mode": "INCREMENTAL",
        "checkpoint_interval_ms": 60000,
    },
    target={"type": "DORIS", "database": "ods"},
)

#4.2 自动映射到 Ontology 对象

Python
platform.ontology.create_mapping(
    source_table="ods.orders",
    object_type="Order",
    field_mappings={
        "order_id":     {"property": "orderId",     "type": "STRING",    "primary_key": True},
        "customer_id":  {"property": "customerId",  "type": "STRING"},
        "product_name": {"property": "productName", "type": "STRING"},
        "total_amount": {"property": "totalAmount",  "type": "DOUBLE"},
        "status":       {"property": "status",       "type": "STRING"},
        "created_at":   {"property": "createdAt",    "type": "TIMESTAMP"},
    },
    sync_mode="REAL_TIME",
    link_rules=[{
        "link_type": "orderedBy",
        "target_object_type": "Customer",
        "source_field": "customer_id",
        "target_field": "customerId",
    }],
)

#4.3 OQL 实时查询

Python
result = platform.oql.execute("""
    SELECT o.orderId, o.productName, o.totalAmount, c.name AS customerName
    FROM Order o
    JOIN o.orderedBy c
    WHERE o.totalAmount > 10000
      AND o.createdAt > NOW() - INTERVAL 1 HOUR
    ORDER BY o.totalAmount DESC
    LIMIT 20
""")

for row in result.rows:
    print(f"Order {row['orderId']}: {row['customerName']} - {row['totalAmount']:.2f}")

#5. 监控与运维

#5.1 关键指标

指标含义告警阈值
numRecordsInPerSecond每秒输入记录数突变 > 200% 或降为 0
numRecordsOutPerSecond每秒输出记录数与输入差距 > 20%
currentFetchEventTimeLagCDC 事件延迟> 60 秒
lastCheckpointDurationCheckpoint 耗时> 30 秒
numberOfFailedCheckpoints失败次数> 0

#5.2 常见问题排查

问题原因解决方案
任务启动失败binlog 已过期增大 expire_logs_days,全量重同步
全量阶段 OOM表太大减小 chunk.size,增加 TM 内存
Doris 写入超时Stream Load 慢增大超时参数,减小批次
Schema 变更中断DDL 处理异常启用 schema.change.behavior: evolve
数据延迟增大Sink 吞吐不足增加并行度
数据重复Checkpoint 不一致启用 Exactly-Once + 2PC

#5.3 灾难恢复

Bash
# 创建 Savepoint
curl -X POST http://flink-jobmanager:8081/jobs/{JOB_ID}/savepoints \
  -d '{"cancel-job": false, "target-directory": "file:///opt/flink/savepoints"}'

# 从 Savepoint 恢复
./bin/flink run \
  -s file:///opt/flink/savepoints/savepoint-xxxxxx \
  -c com.onto.pipeline.cdc.ERPSyncPipeline \
  erp-sync-pipeline.jar

#6. 生产环境最佳实践

#6.1 容量规划

源表规模变更 TPS推荐并行度TM 内存
< 100 万< 10022 GB
100-1000 万100-100044 GB
1000 万-1 亿1000-500088 GB
> 1 亿> 500016+16 GB+

#6.2 上线检查清单

  • MySQL binlog 格式为 ROW,binlog_row_image 为 FULL
  • binlog 保留至少 3 天
  • CDC 用户权限最小化
  • Flink checkpoint 使用持久化存储
  • 监控告警已配置
  • Doris 目标库 Schema 已创建
  • 网络连通性已验证
  • Savepoint 恢复流程已演练

#总结

本教程完整介绍了使用 Flink CDC 实现实时数据同步的全流程:CDC 原理、环境配置、Flink SQL 和 DataStream API 两种实现方式、Ontology 自动映射、运维监控和生产最佳实践。Flink CDC 是 coomia-dip 数据层的"血液循环系统",让企业各业务系统的数据实时汇聚到统一的 Ontology 模型中。

下一篇:[S12-16] Temporal 工作流编排指南 上一篇:[S12-14] gRPC 自定义服务开发指南