一、依赖与环境准备

1.1 Maven 依赖(应用工程)

在自定义 Flink 程序(Java/Scala)中使用 TiDB CDC,只需要在 pom.xml 中引入:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-tidb-cdc</artifactId>
    <version>3.5.0</version>
</dependency>

连接 TiDB 底层实际是访问 TiKV + PD 集群,因此这个 connector 会自动携带 TiDB/TiKV 客户端依赖,你只需要保证:

  • Flink 版本与 3.5.0 对应的兼容矩阵正常
  • 集群节点能访问 PD 地址(内网/外网连通)

1.2 Flink SQL Client / 集群

如果你直接用 Flink SQL 做开发,可以下载:

  • flink-sql-connector-tidb-cdc-3.5.0.jar

拷贝到:

<FLINK_HOME>/lib/

重启 Flink 集群之后,就可以在 SQL 中使用:

'connector' = 'tidb-cdc'

来创建 TiDB CDC 表了。

二、在 Flink SQL 中创建 TiDB CDC 表

先看一个完整示例,直观感受一下:

-- 每 3 秒做一次 checkpoint
SET 'execution.checkpointing.interval' = '3s';

-- 注册 TiDB 中的 orders 表为 CDC 源表
CREATE TABLE orders (
    order_id      INT,
    order_date    TIMESTAMP(3),
    customer_name STRING,
    price         DECIMAL(10, 5),
    product_id    INT,
    order_status  BOOLEAN,
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector'                = 'tidb-cdc',
    'tikv.grpc.timeout_in_ms'  = '20000',         -- TiKV gRPC 超时
    'pd-addresses'             = 'localhost:2379',-- PD 地址列表
    'database-name'            = 'mydb',
    'table-name'               = 'orders'
);

-- 像查普通表一样查:背后其实是快照 + binlog 流
SELECT * FROM orders;

关键点解释:

  1. pd-addresses

    • TiDB 集群的 PD(Placement Driver)地址列表,例如:

      • 单节点:'127.0.0.1:2379'
      • 集群:'10.0.0.1:2379,10.0.0.2:2379,10.0.0.3:2379'
    • Flink CDC 实际是通过 PD 找到 TiKV 节点,进行快照扫描和 binlog 订阅。

  2. database-name / table-name

    • 需要指定具体 database 与 table:

      • database-name = 'mydb'
      • table-name = 'orders'
  3. PRIMARY KEY(order_id) NOT ENFORCED

    • 告诉 Flink:「这张表的语义主键是 order_id
    • 但不会去数据库中校验约束,属于逻辑主键(NOT ENFORCED);
    • 对下游 Upsert Sink(如 upsert-kafka、Hudi/Paimon 主键表)非常关键。
  4. execution.checkpointing.interval 开启 checkpoint,是实现 Exactly-Once 的前提。

三、核心 Connector 参数详解

官方提供了一系列可配置参数,下面挑常用、关键的说明一下。

3.1 基础连接参数

参数名必填说明
connector固定为 tidb-cdc
database-name要监控的 TiDB 数据库名
table-name要监控的表名,例如 orders
pd-addressesTiKV 集群 PD 地址列表,TiDB 的控制面入口

TiDB CDC 与 MySQL CDC 最大的区别是:不再直接连 MySQL 协议端口,而是连 PD+TiKV
因此你不会看到 hostname / port 这样的传统数据库参数,而是接入到 TiKV 集群。

3.2 启动模式:scan.startup.mode

控制 TiDB CDC 第一次启动时,从哪里开始读:

模式默认含义
initial✔️对表做一次结构 + 数据的快照,然后从当前 binlog 继续追数
latest-offset只读取结构快照,不读历史数据,从当前时刻开始接收变更

配置示例:

WITH (
    'connector'        = 'tidb-cdc',
    'pd-addresses'     = 'tidb-pd:2379',
    'database-name'    = 'mydb',
    'table-name'       = 'orders',
    'scan.startup.mode'= 'latest-offset'
)

实战建议:

  • 做全量 + 增量同步到数仓、ODS、Kafka 主题时 → 用 initial
  • 只关心未来新增/更新,比如实时风控规则只看未来订单 → 用 latest-offset

3.3 TiKV 客户端参数

TiDB CDC 底层其实是一个 TiKV Client,直接从 TiKV 的 Region + Raft 日志 拉数据。
因此,你也可以通过一系列 tikv.* 参数微调读性能:

参数名默认说明
tikv.grpc.timeout_in_ms-gRPC 请求超时(毫秒)
tikv.grpc.scan_timeout_in_ms-扫描请求超时
tikv.batch_get_concurrency20并发批量读取的并发度
tikv.*-所有未专门列出的 TiDB/TiKV 客户端参数透传

在网络不稳定、跨 IDC/跨云访问 TiDB 集群时,可以适当调大 timeout,或者降低 batch 并发,避免过多重试。

3.4 内网/外网 IP 映射:host-mapping

如果 TiDB 集群部署在 内网,而 Flink 集群在 外网/K8s 另一套网络空间
这时 TiKV 对外暴露的 IP 和 Flink 能访问到的 IP 可能不一致,就需要配置 host 映射:

host-mapping = "192.168.0.2:8.8.8.8;192.168.0.3:9.9.9.9"

含义:

  • TiKV 节点对 PD 报告的地址是 192.168.0.2/3(内网)
  • Flink 想访问时,必须走 8.8.8.8/9.9.9.9(外网/NAT 出口)
  • connector 会在内部把地址替换之后再连。

这个对「云上 TiDB + 自建 Flink」或者「两个 VPC 之间互通不完善」的场景非常重要。

四、Exactly-Once 与多线程并行读取

4.1 Exactly-Once 语义

TiDB CDC 高度集成了 Flink 的 checkpoint 机制:

  • 快照阶段:按照一致性快照读取 TiKV 中某个 snapshot 的数据;
  • 增量阶段:从 TiKV CDC/Binlog 流中持续消费变更事件;
  • 每次 checkpoint 会持久化消费位点 / 任务状态,故障恢复后可以从正确位置继续。

配合:

  • execution.checkpointing.interval
  • 合理的重启策略

常规场景可以提供 Exactly-Once 语义(对于无主键 / 下游幂等写的情况,语义会弱化为 At-Least-Once + 人工幂等)。

4.2 多线程并行读取(TiDB 的一大优势)

对比 MySQL/Oracle/SQLServer 这些传统单机数据库,TiDB 最大的特点是:

TiDB CDC 源可以多线程 / 多 task 并行读取

因为 TiKV 本身是分 Region、分 Store 的分布式集群,可以:

  • 将 Table 分为多个范围(类似分片);
  • 不同 Flink Source 子任务并行读取不同 Region;
  • 快照阶段和 binlog 流阶段都可以做到更好的扩展性。

在 SQL 层你不需要感知,只要增加 Source 并行度即可:

-- Flink SQL CLI 中可以通过 table.exec.resource.default-parallelism 控制默认并行度
SET 'table.exec.resource.default-parallelism' = '4';

在 DataStream API 中:

env.fromSource(tidbSource, ..., "TiDBSource")
   .setParallelism(4);

五、DataStream API 使用示例

如果你需要对 TiDB 原始事件做更底层处理,或者想把 CDC 事件以自定义格式输出,可以直接使用 TiDBSource

5.1 完整示例:自定义快照/变更反序列化

import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.util.Collector;

import org.apache.flink.cdc.connectors.tidb.TDBSourceOptions;
import org.apache.flink.cdc.connectors.tidb.TiDBSource;
import org.apache.flink.cdc.connectors.tidb.TiKVChangeEventDeserializationSchema;
import org.apache.flink.cdc.connectors.tidb.TiKVSnapshotEventDeserializationSchema;

import org.tikv.kvproto.Cdcpb;
import org.tikv.kvproto.Kvrpcpb;

import java.util.HashMap;

public class TiDBSourceExample {

    public static void main(String[] args) throws Exception {

        SourceFunction<String> tidbSource =
            TiDBSource.<String>builder()
                .database("mydb")           // 捕获的数据库
                .tableName("products")      // 捕获的表
                .tiConf(
                    TDBSourceOptions.getTiConfiguration(
                        "localhost:2399", new HashMap<>()))
                // 快照阶段反序列化
                .snapshotEventDeserializer(
                    new TiKVSnapshotEventDeserializationSchema<String>() {
                        @Override
                        public void deserialize(
                                Kvrpcpb.KvPair record,
                                Collector<String> out) {
                            out.collect("SNAPSHOT: " + record.toString());
                        }

                        @Override
                        public TypeInformation<String> getProducedType() {
                            return BasicTypeInfo.STRING_TYPE_INFO;
                        }
                    })
                // 增量变更阶段反序列化
                .changeEventDeserializer(
                    new TiKVChangeEventDeserializationSchema<String>() {
                        @Override
                        public void deserialize(
                                Cdcpb.Event.Row record,
                                Collector<String> out) {
                            out.collect("CHANGE: " + record.toString());
                        }

                        @Override
                        public TypeInformation<String> getProducedType() {
                            return BasicTypeInfo.STRING_TYPE_INFO;
                        }
                    })
                .build();

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(3000);

        env.addSource(tidbSource)
           .print()
           .setParallelism(1);

        env.execute("Print TiDB Snapshot + Binlog");
    }
}

这里有几个关键点:

  1. TiKVSnapshotEventDeserializationSchema

    • 只在快照阶段调用,输入类型为 Kvrpcpb.KvPair
    • 你可以在内部把 Key/Value 解码成 JSON 或 POJO,再往 Collector 里发射。
  2. TiKVChangeEventDeserializationSchema

    • 在增量阶段调用,输入类型为 Cdcpb.Event.Row
    • 事件里包含 Key、Column 值、操作类型(插入/更新/删除)等,可以做更细致的处理。
  3. TDBSourceOptions.getTiConfiguration

    • 这里需要传 PD 地址(示例是 "localhost:2399",真实环境下换成你的 PD 地址);
    • 以及其他 TiKV Client 配置(HashMap 中可以继续放 tiConf 参数)。

实战中你通常会:

  • 把 snapshot + change 都解成统一结构(比如 JSON),
  • 后续统一用 Flink DataStream 转成动态表,或直接写出到 Kafka / ES / OLAP。

六、元数据列:库名、表名、变更时间戳

TiDB CDC 也支持和其他 CDC connector 一样,暴露一些只读元数据列:

key类型含义
table_nameSTRING NOT NULL当前记录所属表名
database_nameSTRING NOT NULL当前记录所属数据库名
op_tsTIMESTAMP_LTZ(3) NOT NULL变更在 TiDB 中发生的时间;快照数据为 0

扩展 DDL 示例:

CREATE TABLE products (
    db_name      STRING METADATA FROM 'database_name' VIRTUAL,
    table_name   STRING METADATA FROM 'table_name'   VIRTUAL,
    operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,

    order_id      INT,
    order_date    TIMESTAMP(0),
    customer_name STRING,
    price         DECIMAL(10, 5),
    product_id    INT,
    order_status  BOOLEAN,
    PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
    'connector'               = 'tidb-cdc',
    'tikv.grpc.timeout_in_ms' = '20000',
    'pd-addresses'            = 'localhost:2379',
    'database-name'           = 'mydb',
    'table-name'              = 'orders'
);

有了这些字段,你可以:

  • 按库/表分区写入下游文件系统;
  • operation_ts 做窗口统计或迟到数据处理;
  • 在审计日志中记录「这条数据来自哪张表」。

七、TiDB → Flink SQL 类型映射

TiDB 类型与 Flink SQL 类型的映射总体比较自然,但有几个地方需要特别注意。

7.1 数值类型

TiDB 类型Flink SQL 类型说明
TINYINTTINYINT
SMALLINT, TINYINT UNSIGNEDSMALLINT
INT, MEDIUMINT, SMALLINT UNSIGNEDINT
BIGINT, INT UNSIGNEDBIGINT
BIGINT UNSIGNEDDECIMAL(20, 0)避免溢出
FLOATFLOAT
REAL, DOUBLEDOUBLE
NUMERIC(p, s), DECIMAL(p, s) 且 p ≤ 38DECIMAL(p, s)正常映射
NUMERIC/DECIMAL p > 38STRINGTiDB 支持精度到 65,而 Flink DECIMAL 只支持到 38,需要映射成 STRING 避免精度损失

实战提示:
金额类字段一般 DECIMAL(18,2) 足够,如果你确实在 TiDB 用了精度 > 38 的 DECIMAL,
记得在 Flink 侧用 STRING 接一下,再在下游做高精度处理。

7.2 布尔与时间类型

TiDB 类型Flink SQL 类型
BOOLEAN / TINYINT(1) / BIT(1)BOOLEAN
DATEDATE
TIME§TIME§
TIMESTAMP§TIMESTAMP_LTZ§
DATETIME§TIMESTAMP§

这里最大的区别是:

  • TIMESTAMPTIMESTAMP_LTZ:带时区语义,适合跨时区的业务;
  • DATETIME 则更像「本地时间」,没有时区信息,对应普通 TIMESTAMP

7.3 字符与二进制类型

TiDB 类型Flink SQL 类型
CHAR(n)CHAR(n)
VARCHAR(n)VARCHAR(n)
TINYTEXT/TEXT/MEDIUMTEXT/LONGTEXTSTRING
BIT(n)BINARY(ceil(n/8))
BINARY(n)BINARY(n)
TINYBLOB/BLOB/MEDIUMBLOB/LONGBLOBBYTES

限制:
BLOB 类型目前只支持 长度 ≤ 2,147,483,647 (2^31 - 1) 的数据。

7.4 复杂类型:ENUM / JSON / SET

TiDB 类型Flink SQL 类型说明
YEARINT
ENUMSTRING
JSONSTRINGJSON 文本串,后续可再解析
SETARRAYTiDB SET 是多选字符串集合,映射成字符串数组

JSON/SET 在下游处理时往往会配合:

  • Flink SQL 的 JSON 函数 / UDF;
  • 或在 DataStream 中直接用 Jackson / Gson 解析。

八、实战落地建议小结

最后,用一份 checklist 帮你快速回顾 TiDB CDC 落地时需要关注的点:

  1. 集群连通与 PD 地址

    • 确认 Flink TaskManager 能访问 PD 地址;
    • 如存在内外网问题,配置好 host-mapping
  2. checkpoint 与容错策略

    • 打开 checkpoint(例如 30s/1min 一次);
    • 配合合理的 restart-strategy。
  3. 启动模式选择

    • 首次全量 + 增量同步 → scan.startup.mode='initial'
    • 只关心未来变更 → scan.startup.mode='latest-offset'
  4. 并行度与资源

    • 利用 TiDB CDC 的多线程优势,合理提高 Source 并行度;
    • 注意 TiKV/网络层的负载上限,避免打挂生产集群。
  5. Decimal 高精度字段

    • p > 38 的 DECIMAL 统一按 STRING 映射;
    • 在下游系统中做高精度运算或只用于展示。
  6. 复杂类型处理

    • JSON → STRING(再解析);
    • SET → ARRAY(多值字段)。

做到这些,基本就可以把 TiDB 里的核心业务表顺利接到 Flink 实时链路上,
无论是做 TiDB → Kafka → 实时数仓,还是 TiDB → Lakehouse(Iceberg/Hudi/Paimon)
TiDB CDC Connector 都是一块非常关键的拼图。

更多推荐