用 Flink TiDB CDC 打通 TiDB 的实时数据通道
一、依赖与环境准备
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;
关键点解释:
-
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 订阅。
-
-
database-name/table-name-
需要指定具体 database 与 table:
database-name = 'mydb'table-name = 'orders'
-
-
PRIMARY KEY(order_id) NOT ENFORCED- 告诉 Flink:「这张表的语义主键是
order_id」 - 但不会去数据库中校验约束,属于逻辑主键(NOT ENFORCED);
- 对下游 Upsert Sink(如 upsert-kafka、Hudi/Paimon 主键表)非常关键。
- 告诉 Flink:「这张表的语义主键是
-
给
execution.checkpointing.interval开启 checkpoint,是实现 Exactly-Once 的前提。
三、核心 Connector 参数详解
官方提供了一系列可配置参数,下面挑常用、关键的说明一下。
3.1 基础连接参数
| 参数名 | 必填 | 说明 |
|---|---|---|
| connector | 是 | 固定为 tidb-cdc |
| database-name | 是 | 要监控的 TiDB 数据库名 |
| table-name | 是 | 要监控的表名,例如 orders |
| pd-addresses | 是 | TiKV 集群 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_concurrency | 20 | 并发批量读取的并发度 |
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");
}
}
这里有几个关键点:
-
TiKVSnapshotEventDeserializationSchema- 只在快照阶段调用,输入类型为
Kvrpcpb.KvPair; - 你可以在内部把 Key/Value 解码成 JSON 或 POJO,再往
Collector里发射。
- 只在快照阶段调用,输入类型为
-
TiKVChangeEventDeserializationSchema- 在增量阶段调用,输入类型为
Cdcpb.Event.Row; - 事件里包含 Key、Column 值、操作类型(插入/更新/删除)等,可以做更细致的处理。
- 在增量阶段调用,输入类型为
-
TDBSourceOptions.getTiConfiguration- 这里需要传 PD 地址(示例是
"localhost:2399",真实环境下换成你的 PD 地址); - 以及其他 TiKV Client 配置(
HashMap中可以继续放 tiConf 参数)。
- 这里需要传 PD 地址(示例是
实战中你通常会:
- 把 snapshot + change 都解成统一结构(比如 JSON),
- 后续统一用 Flink DataStream 转成动态表,或直接写出到 Kafka / ES / OLAP。
六、元数据列:库名、表名、变更时间戳
TiDB CDC 也支持和其他 CDC connector 一样,暴露一些只读元数据列:
| key | 类型 | 含义 |
|---|---|---|
| table_name | STRING NOT NULL | 当前记录所属表名 |
| database_name | STRING NOT NULL | 当前记录所属数据库名 |
| op_ts | TIMESTAMP_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 类型 | 说明 |
|---|---|---|
| TINYINT | TINYINT | |
| SMALLINT, TINYINT UNSIGNED | SMALLINT | |
| INT, MEDIUMINT, SMALLINT UNSIGNED | INT | |
| BIGINT, INT UNSIGNED | BIGINT | |
| BIGINT UNSIGNED | DECIMAL(20, 0) | 避免溢出 |
| FLOAT | FLOAT | |
| REAL, DOUBLE | DOUBLE | |
| NUMERIC(p, s), DECIMAL(p, s) 且 p ≤ 38 | DECIMAL(p, s) | 正常映射 |
| NUMERIC/DECIMAL p > 38 | STRING | TiDB 支持精度到 65,而 Flink DECIMAL 只支持到 38,需要映射成 STRING 避免精度损失 |
实战提示:
金额类字段一般DECIMAL(18,2)足够,如果你确实在 TiDB 用了精度 > 38 的 DECIMAL,
记得在 Flink 侧用 STRING 接一下,再在下游做高精度处理。
7.2 布尔与时间类型
| TiDB 类型 | Flink SQL 类型 |
|---|---|
| BOOLEAN / TINYINT(1) / BIT(1) | BOOLEAN |
| DATE | DATE |
| TIME§ | TIME§ |
| TIMESTAMP§ | TIMESTAMP_LTZ§ |
| DATETIME§ | TIMESTAMP§ |
这里最大的区别是:
TIMESTAMP→TIMESTAMP_LTZ:带时区语义,适合跨时区的业务;DATETIME则更像「本地时间」,没有时区信息,对应普通TIMESTAMP。
7.3 字符与二进制类型
| TiDB 类型 | Flink SQL 类型 |
|---|---|
| CHAR(n) | CHAR(n) |
| VARCHAR(n) | VARCHAR(n) |
| TINYTEXT/TEXT/MEDIUMTEXT/LONGTEXT | STRING |
| BIT(n) | BINARY(ceil(n/8)) |
| BINARY(n) | BINARY(n) |
| TINYBLOB/BLOB/MEDIUMBLOB/LONGBLOB | BYTES |
限制:
BLOB 类型目前只支持 长度 ≤ 2,147,483,647 (2^31 - 1) 的数据。
7.4 复杂类型:ENUM / JSON / SET
| TiDB 类型 | Flink SQL 类型 | 说明 |
|---|---|---|
| YEAR | INT | |
| ENUM | STRING | |
| JSON | STRING | JSON 文本串,后续可再解析 |
| SET | ARRAY | TiDB SET 是多选字符串集合,映射成字符串数组 |
JSON/SET 在下游处理时往往会配合:
- Flink SQL 的 JSON 函数 / UDF;
- 或在 DataStream 中直接用 Jackson / Gson 解析。
八、实战落地建议小结
最后,用一份 checklist 帮你快速回顾 TiDB CDC 落地时需要关注的点:
-
集群连通与 PD 地址
- 确认 Flink TaskManager 能访问 PD 地址;
- 如存在内外网问题,配置好
host-mapping。
-
checkpoint 与容错策略
- 打开 checkpoint(例如 30s/1min 一次);
- 配合合理的 restart-strategy。
-
启动模式选择
- 首次全量 + 增量同步 →
scan.startup.mode='initial'; - 只关心未来变更 →
scan.startup.mode='latest-offset'。
- 首次全量 + 增量同步 →
-
并行度与资源
- 利用 TiDB CDC 的多线程优势,合理提高 Source 并行度;
- 注意 TiKV/网络层的负载上限,避免打挂生产集群。
-
Decimal 高精度字段
- p > 38 的 DECIMAL 统一按 STRING 映射;
- 在下游系统中做高精度运算或只用于展示。
-
复杂类型处理
- JSON → STRING(再解析);
- SET → ARRAY(多值字段)。
做到这些,基本就可以把 TiDB 里的核心业务表顺利接到 Flink 实时链路上,
无论是做 TiDB → Kafka → 实时数仓,还是 TiDB → Lakehouse(Iceberg/Hudi/Paimon),
TiDB CDC Connector 都是一块非常关键的拼图。
更多推荐
所有评论(0)