Day 10|Doris 流式同步与 CDC:Routine Load、Flink Connector 与 Streaming Job

Day 9 解决了批量数据怎样进入 Doris:客户端可以通过 Stream Load 推送文件,Doris 可以通过 TVF、Catalog 或历史上的 Broker Load 读取远端数据,内部表之间也可以使用 INSERT INTO ... SELECT 完成加工。批量任务拥有清晰的开始和结束,一次提交对应一个相对明确的数据集合。
流式同步把问题延伸到了时间轴上。Kafka 会持续产生消息,MySQL Binlog 和 PostgreSQL WAL 会不断记录变化,业务表可能新增字段,网络会抖动,任务会重启,数据会重放,旧事件也可能晚于新事件抵达。系统需要长期保存读取进度,并在每次恢复时找到正确的继续位置。这里的“实时”同时包含四层含义:数据持续进入、状态持续推进、故障后能够恢复、查询结果持续逼近源端真实状态。
因此,流式链路的设计重点不止是“每秒能写多少行”。一套可以进入生产的方案,还要回答下面这些问题:
- Kafka Offset、Flink Checkpoint 或数据库日志位点由谁维护;
- 一批数据写入成功时,读取进度是否同步提交;
- 任务在提交前宕机、提交后丢失响应时会发生什么;
INSERT、UPDATE、DELETE怎样映射到 Doris 表模型;- 旧事件晚到时,怎样防止它覆盖新状态;
- 源表加列、改类型或改主键时,链路怎样处理;
- 任务恢复后,怎样证明无缺失、无异常重复、状态无回退;
- 延迟升高时,问题位于源端、流处理、Doris 写入,还是后台 Compaction。
本篇将围绕 Routine Load、Doris Kafka Connector、Flink Doris Connector、Flink CDC 和 Doris Streaming Job 建立一张完整地图。重点放在原理、语义、选型和生产验收。完成本篇后,学习者应当能够独立设计一条 Kafka 或数据库 CDC 到 Doris 的实时同步链路,并写出相应的运行手册与验收标准。
本日学习目标
完成本篇后,你应当能够:
- 根据数据源、格式、处理复杂度和一致性要求选择实时同步方式;
- 解释 Routine Load 的 Job、Task、Kafka Partition 和 Doris 事务之间的关系;
- 说明 Routine Load 如何通过消费进度持久化和提交校验实现 Exactly-Once;
- 正确设置起始 Offset、任务并发和微批阈值;
- 使用
SHOW ROUTINE LOAD判断任务状态、消费进度和暂停原因; - 解释 Flink Checkpoint 与 Stream Load 2PC 怎样共同控制数据可见性;
- 使用 Unique Key、Delete Sign 和 Sequence Column 承接 CDC 事件;
- 理解 Flink CDC 整库同步、自动建表和 Schema 演进的边界;
- 区分 Streaming Job 的 SQL Mapping Sync 与 Auto Table Creation Sync;
- 建立覆盖 Lag、吞吐、质量、恢复和对账的流式任务验收体系。
版本口径
截至 2026-08-16,Apache Doris 下载页将 4.1.3 标记为 Latest,将 4.0.8 标记为 Stable。本文以 4.x 通用能力为主。Streaming Job、CDC Stream TVF 以及部分 Connector 新参数需要在目标小版本上做语法与行为验证。Flink Doris Connector 当前仓库最新发布为 26.1.1,官方开发文档列出的兼容范围覆盖 Flink 1.15–1.20 和 2.0–2.2;Flink 1.x 与 2.x 对 Java 版本的要求不同,部署前应按照实际组合核验。
实验前提
Routine Load 实验需要 Kafka。Flink 实验需要独立 Flink 集群和匹配版本的 Connector Jar。Streaming Job 的数据库 CDC 实验需要 MySQL Binlog 或 PostgreSQL Logical Replication。完整包提供模板、样例消息和验收命令,当前制作环境没有连接实际 Kafka、Flink、MySQL、PostgreSQL 和 Doris 集群,运行结果需要在 Day 5 的目标环境中完成最终确认。
一、先建立流式同步坐标系

实时链路的产品名称很多,判断方法可以保持简单。先回答五个问题,再进入参数层。
1. 数据从哪里来
常见来源包括:
- Kafka Topic 中的日志、订单事件、埋点和 IoT 数据;
- MySQL Binlog、PostgreSQL WAL、Oracle Redo 等数据库变化日志;
- Flink 已经加工完成的实时宽表或聚合流;
- S3、OSS、COS 等对象存储目录中持续产生的新文件。
来源决定读取协议和进度形态。Kafka 使用 Partition Offset,Flink 使用 Checkpoint 和 Operator State,数据库 CDC 使用 Binlog Position、GTID、LSN 或等价日志位点,对象存储持续导入通常使用文件名或已处理文件清单作为进度。
2. 数据格式是什么
Routine Load 直接消费 Kafka 时支持 CSV 和 JSON。Avro、Protobuf、Debezium 等格式通常交给 Doris Kafka Connector 或上游 Flink 解析。数据库 CDC 事件还包含操作类型、事务边界、源表信息和日志位点,单纯把事件正文当作普通 JSON 处理,容易丢失删除和顺序语义。
3. 是否需要实时计算
直写链路只做字段映射、类型转换和简单过滤时,Routine Load 或 Streaming Job 已经足够。需要跨流 Join、窗口聚合、维表关联、复杂状态计算时,Flink 更合适。计算逻辑越复杂,状态管理、Checkpoint 时间和反压传播越需要纳入容量规划。
4. 谁维护进度
这是流式系统最重要的分界线之一:
| 方案 | 进度与状态管理者 | 典型进度 |
|---|---|---|
| Routine Load | Doris FE | Kafka Partition Offset |
| Doris Kafka Connector | Kafka Connect | Connector Offset |
| Flink Doris Connector | Flink | Checkpoint / Savepoint |
| Flink CDC | Flink CDC Source | Snapshot State + Binlog/WAL Position |
| Doris Streaming Job | Doris Job Framework | CDC Offset 或文件进度 |
进度管理者承担恢复责任。任务重启时,它需要知道最后一次已经形成有效下游状态的位置,并从该位置继续。
5. 需要哪种一致性语义
“Exactly-Once”必须注明范围。Routine Load 的语义覆盖 Kafka 消息消费到 Doris 事务可见。Flink Connector 的流式模式通过 Checkpoint 和 Stream Load 2PC 对齐 Flink State 与 Doris 版本。Streaming Job 的 SQL Mapping Sync 支持 Exactly-Once;自动建表整库镜像目前保证 At-Least-Once。At-Least-Once 可以通过 Unique Key 主键去重降低重复影响,仍需关注顺序、删除和无主键表。
这五个问题回答清楚以后,选型会变得直接:
- Kafka CSV/JSON 直写,优先评估 Routine Load;
- Kafka Avro、Protobuf 或 Debezium,评估 Doris Kafka Connector;
- 企业已经建设 Flink 平台,实时加工逻辑较多,使用 Flink Doris Connector;
- 数据库整库、多表、自动建表和 DDL 同步,使用 Flink CDC + Doris Connector;
- 希望减少外部组件,单表需要字段映射与过滤,评估 Streaming Job SQL Mapping;
- 希望镜像一组 MySQL/PostgreSQL 表,且无需转换,评估 Streaming Job Auto Table Creation。
二、Routine Load:Doris 原生的 Kafka 持续导入任务

你提供的培训材料用一张图概括了 Routine Load:FE 负责 Job 调度,把 Kafka Partition 动态分配给 BE Task;BE 拉取消息并写入目标表,消费进度随任务提交。这个框架仍然成立。当前 4.x 官方文档进一步把它拆成两层对象:长期存在的 Routine Load Job,以及不断生成和结束的 Routine Load Task。
1. Job 是长期控制面
执行 CREATE ROUTINE LOAD 后,FE 创建一个长期 Job。Job 保存:
- Kafka Broker、Topic 和认证配置;
- 起始 Offset 与当前消费进度;
- 目标数据库、目标表和导入规则;
- 期望并发、微批阈值和质量参数;
- 当前状态、暂停原因、错误信息和统计数据。
Job Scheduler 负责生命周期管理、Task 拆分、失败恢复和状态转换。它不会把整个 Topic 当作一个永不结束的单次事务。
2. Task 是短生命周期执行单元
Job 会按照 Kafka Partition 和并发配置拆成多个 Task。每个 Task:
- 负责一组 Kafka Partition;
- 从已记录 Offset 开始读取;
- 在达到时间、行数或数据量阈值时结束当前微批;
- 把数据写入 Doris;
- 以独立事务提交;
- 提交成功后更新消费进度;
- 由 Job 继续生成下一轮 Task。
这种设计让长期数据流被转换成一系列有边界的微批事务。每个事务规模可控,失败时也能从明确的 Offset 重试。
3. Task 并发怎样确定
实际并发由三个上限共同决定:
actual_task_num = min(
kafka_topic_partition_num,
desired_concurrent_number,
max_routine_load_task_concurrent_num
)
Kafka Topic 只有 6 个 Partition 时,把 desired_concurrent_number 设置为 32 也无法得到 32 个有效消费 Task。扩展 Routine Load 吞吐时,需要同时检查 Kafka Partition 数、FE 的集群级并发上限、BE 数量和目标表 Tablet 分布。
4. 支持的数据格式与处理能力
Routine Load 当前直接支持 CSV 和 JSON。它可以完成:
- 列映射;
- 派生列计算;
WHERE过滤;- JSONPath 和 JSON Root 提取;
- 指定目标分区;
- Unique Key 更新与删除;
- 一条 Topic 动态路由到多张 Doris 表。
一流多表模式要求 Kafka 消息携带目标表名,表名需要与 Doris 表匹配,同时存在列映射等限制。它适合统一采集 Topic 中已经带有明确表路由信息的场景。复杂业务转换继续交给 Flink 或独立处理层,能够降低导入链路的调试成本。
5. Routine Load 与 Doris Kafka Connector 的边界
Routine Load 直接在 Doris 中消费 Kafka,部署面最小,CSV 和 JSON 场景优先采用这条路径。Doris Kafka Connector 运行在 Kafka Connect Framework 中,由 Kafka Connect 管理 Worker、Task、Offset 和故障转移。它适合以下需求:
- Topic 使用 Avro 或 Protobuf;
- 上游采用 Schema Registry;
- 已经使用 Debezium 把数据库变化写入 Kafka;
- 团队希望把采集任务统一纳入 Kafka Connect 管理;
- 需要依靠 Kafka Connect Worker 水平扩展。
两条链路都需要关注目标表模型和删除语义。Debezium 消息包含 before、after、op 和源端位点等结构,进入 Doris 前要明确主键、事件操作类型、字段映射与顺序规则。使用 Kafka Connector 并不会自动消除业务乱序,Sequence Column 和源端版本仍要纳入表设计。
选型时还要考虑运维归属。Routine Load 的运行状态和消费进度在 Doris 中查看;Kafka Connector 的 Worker、Connector、Task 和 Offset 需要在 Kafka Connect 平台中监控。企业已经有成熟 Kafka Connect 平台时,后者可以复用现有治理;团队希望减少中间组件时,Routine Load 更轻。
三、Routine Load 怎样实现 Exactly-Once

Exactly-Once 容易被理解成“不会发生任何重试”。真实系统会重试,关键在于重试不会形成第二次有效写入。Routine Load 依靠两组机制完成这一点。
1. 消费进度与事务结果一起持久化
每个 Task 读取一段 Kafka Offset,例如 Partition 0 的 [100, 200)。它把这段消息写入 Doris,事务提交成功时,同时把新的消费进度写入 FE Edit Log,并同步到 FE Follower。恢复后,新 Task 从 Offset 200 继续。
数据版本和消费进度需要保持一致:
- 数据未提交,Offset 不能推进;
- Offset 已推进,对应数据必须已经形成有效版本;
- 查询只能读到已经发布的可见版本。
2. 提交前检查 Task 是否仍然有效
Task 提交时,FE 会检查它是否仍存在于当前有效任务列表。旧 Task 超时后被重新调度,或 Job 状态已经变化时,迟到提交会被拒绝。这一步用于避免两个 Task 对同一段 Offset 同时形成有效提交。
3. 两类典型故障
场景 A:BE 写完数据,在 Commit 前宕机。
事务尚未提交,Rowset 对查询不可见,Offset 也没有推进。恢复后,新的 Task 重新读取同一段 Kafka 消息。旧的临时数据由系统清理。
场景 B:事务已经提交,客户端或 Task 没有收到成功响应。
FE 已经记录新版本和 Offset。恢复后,新 Task 从更新后的 Offset 继续。旧 Task 再次提交时会被有效性校验拒绝。
4. 语义边界需要写进方案
Routine Load 的 Exactly-Once 解决 Kafka 到 Doris 的一次有效消费。业务侧仍需处理:
- 上游在 Kafka 中本身就产生了重复业务事件;
- 相同业务主键的更新事件发生乱序;
- 业务事件缺少稳定主键;
- 多 Topic 或多表之间存在跨流事务要求;
- 下游业务同时读取 Doris 和其他系统,需要跨系统一致快照。
这些问题分别需要业务事件 ID、Unique Key、Sequence Column、事务编排或批次水位设计。导入语义和业务语义应分开描述。
四、Routine Load 从建表到运行
下面使用订单状态流演示一条最小链路。目标表使用 Unique Key MoW,以 order_id 作为业务主键,并使用 source_commit_ts 作为顺序列。
CREATE DATABASE IF NOT EXISTS day10_lab;
CREATE TABLE day10_lab.order_state (
order_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
status VARCHAR(32) NOT NULL,
amount DECIMAL(18,2) NOT NULL,
update_time DATETIME(3) NOT NULL,
source_commit_ts BIGINT NOT NULL
)
UNIQUE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 8
PROPERTIES (
"replication_num" = "1",
"enable_unique_key_merge_on_write" = "true",
"function_column.sequence_col" = "source_commit_ts"
);
学习环境只有一个 BE,因此示例副本数为 1。生产环境需要按照高可用要求设置副本。
1. 创建 JSON Routine Load
Kafka 中每条消息为一条 JSON:
{"order_id":1001,"user_id":7,"status":"PAID","amount":199.00,
"update_time":"2026-08-16 10:00:00.000","source_commit_ts":1755338400000}
创建任务:
CREATE ROUTINE LOAD day10_lab.rl_order_state
ON order_state
COLUMNS(order_id, user_id, status, amount, update_time, source_commit_ts)
PROPERTIES (
"format" = "json",
"desired_concurrent_number" = "3",
"max_batch_interval" = "10",
"max_batch_rows" = "200000",
"max_batch_size" = "104857600",
"strict_mode" = "true",
"max_filter_ratio" = "0",
"max_error_number" = "0"
)
FROM KAFKA (
"kafka_broker_list" = "kafka1:9092,kafka2:9092,kafka3:9092",
"kafka_topic" = "order_state_json",
"property.group.id" = "doris_day10_order_state",
"property.kafka_default_offsets" = "OFFSET_BEGINNING"
);
这组阈值用于实验观察,生产环境需要根据消息速率和 SLA 调整。max_batch_interval、max_batch_rows 和 max_batch_size 中任意一个达到,当前 Task 就会结束并提交事务。
2. 查看任务状态
SHOW ROUTINE LOAD FOR day10_lab.rl_order_state\G
SHOW ROUTINE LOAD TASK WHERE JobName = 'rl_order_state'\G
重点字段包括:
State:任务当前状态;CurrentTaskNum:当前 Task 数;Statistic:累计行数、错误行、接收字节等;Progress:各 Kafka Partition 当前 Offset;Lag或相关消费落后信息;ReasonOfStateChanged:状态变化原因;OtherMsg:补充错误信息。
不同 4.x 小版本展示字段可能存在差异,运维平台应以目标版本真实输出为准。
3. 暂停、恢复、调整与停止
PAUSE ROUTINE LOAD FOR day10_lab.rl_order_state;
ALTER ROUTINE LOAD FOR day10_lab.rl_order_state
PROPERTIES (
"desired_concurrent_number" = "6",
"max_batch_interval" = "30"
);
RESUME ROUTINE LOAD FOR day10_lab.rl_order_state;
-- 终止后无法恢复,实验结束前不要执行
-- STOP ROUTINE LOAD FOR day10_lab.rl_order_state;
修改 Routine Load 前必须先 Pause,完成 ALTER ROUTINE LOAD 后再 Resume。PAUSED 状态可以恢复。STOPPED 是终态,停止后需要重新创建任务。变更 Topic、认证或 Offset 前,先记录当前进度和对账水位,避免恢复点失去审计依据。
五、Offset、状态机与性能调节

Routine Load 的主要状态包括:
| 状态 | 含义 | 后续动作 |
|---|---|---|
NEED_SCHEDULE |
等待初次调度或重新调度 | 调度成功后进入 RUNNING |
RUNNING |
正常持续消费 | 可能暂停或停止 |
PAUSED |
任务暂停,未终止 | 手工 Resume 或满足条件后自动恢复 |
STOPPED |
手工永久停止 | 无法恢复 |
CANCELLED |
表、库被删除或出现不可恢复错误 | 无法恢复 |
Kafka 暂时不可用、网络抖动等异常通常能够触发自动恢复。手工执行 PAUSE、数据质量问题和不可恢复错误一般需要人工处理。恢复动作前应先检查 ReasonOfStateChanged,直接反复 Resume 可能让任务在错误条件未消除时持续抖动。

1. 起始 Offset
OFFSET_BEGINNING:从 Kafka 当前仍保留的最早消息开始;OFFSET_END:从当前末尾开始,只接收后续新增消息;- 指定 Partition Offset 或时间:用于精确回放、切换和补数。
Kafka Retention 会清理旧消息。任务暂停时间超过保留期后,原 Offset 可能进入 out of range。生产设计需要让 Kafka 日志保留时间覆盖故障发现、审批、修复和重新同步窗口,并保留源库或归档文件作为最终补数来源。
2. 低延迟与高吞吐
官方最佳实践给出的方向很清晰:
- 低延迟场景可以减小默认 60 秒的
max_batch_interval; - 数据量小、资源敏感时可以降低
desired_concurrent_number; - 高吞吐场景可以把
max_batch_interval增加到 120–180 秒,让每个事务承载更多数据。
每次提交都会产生事务、版本、Rowset 和 Segment。间隔过短会增加 FE 元数据和 Compaction 压力。间隔过长会提高可见延迟,并扩大失败重试批次。合适值取决于消息速率、目标 SLA、分桶数量、索引成本和集群资源。
3. Lag 升高时的排查顺序
- 确认 Job 处于
RUNNING; - 对比 Kafka Partition 数和实际 Task 数;
- 检查
desired_concurrent_number与 FE 上限; - 检查 BE CPU、内存、磁盘和网络;
- 检查过滤率、解析错误和单行超大消息;
- 检查目标表索引、Unique Key 更新和 Delete Bitmap 成本;
- 检查 Rowset、Segment 与 Compaction Score;
- 检查 Kafka Broker 吞吐和跨机房网络。
只提高并发可能把瓶颈从消费端移动到 BE Flush 或 Compaction。调优需要同时观察端到端延迟和存储健康度。
六、Flink Doris Connector:把 Flink State 与 Doris 版本对齐

Flink Doris Connector 适合已经拥有 Flink 实时计算平台,或需要窗口、Join、维表关联和复杂状态计算的团队。Connector 会在 Flink 内存中积累数据,再通过 Stream Load 批量写入 Doris。
1. 两种写入模式
官方当前文档将写入分为两类:
| 模式 | 触发条件 | 一致性 | 适用方向 |
|---|---|---|---|
| Streaming Write | Flink Checkpoint | Exactly-Once | 关键实时 ETL、CDC |
| Batch Write | 时间和数据量阈值 | At-Least-Once | 吞吐优先,可使用主键去重 |
Streaming Write 是默认方式。Flink 启动 Checkpoint 时,Doris Sink 对当前批次执行 Stream Load 2PC Prepare;当整个 Checkpoint 成功,Connector 提交 Doris 事务,版本进入可见状态。Checkpoint 失败时,对应事务不会成为有效版本。
2. sink.label-prefix 为什么必须稳定且全局唯一
2PC 场景下,Connector 会基于 Label Prefix、Subtask 和 Checkpoint 构造事务 Label。Prefix 冲突会造成不同 Flink Job 争用 Label。任务恢复时应从最新 Checkpoint 或 Savepoint 启动,随意丢弃状态并复用原 Prefix,可能出现 Label 已使用或事务状态不匹配。
3. Flink SQL 示例
SET 'execution.checkpointing.interval' = '10s';
CREATE TABLE cdc_mysql_source (
id BIGINT,
name STRING,
status STRING,
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-host',
'port' = '3306',
'username' = 'cdc_user',
'password' = '${MYSQL_PASSWORD}',
'database-name' = 'biz_db',
'table-name' = 'user_state'
);
CREATE TABLE doris_sink (
id BIGINT,
name STRING,
status STRING,
update_time TIMESTAMP(3)
) WITH (
'connector' = 'doris',
'fenodes' = 'fe1:8030,fe2:8030',
'table.identifier' = 'day10_lab.user_state',
'username' = 'flink_writer',
'password' = '${DORIS_PASSWORD}',
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true',
'sink.enable-delete' = 'true',
'sink.enable-2pc' = 'true',
'sink.label-prefix' = 'day10_user_state_v1'
);
INSERT INTO doris_sink
SELECT id, name, status, update_time
FROM cdc_mysql_source;
sink.enable-delete=true 会把 Flink RowKind 中的删除事件映射为 Doris 删除语义,目标表需要使用 Unique Key。启用 Batch Mode 后,写入时机脱离 Checkpoint,Exactly-Once 不再由 Connector 保证。生产配置需要显式写出模式,避免默认值在升级或模板复用时产生理解偏差。
4. Checkpoint 间隔怎样影响延迟
Checkpoint 间隔越短,正常情况下可见延迟越低,Checkpoint 协调、状态快照和 Doris 事务频率也会增加。状态较大的 Flink 作业还需要关注 Checkpoint Duration、Alignment Time、失败率和存储带宽。实时 SLA 应分解为:
源端产生延迟
+ Flink 处理与反压
+ 等待 Checkpoint
+ Stream Load 写入
+ Doris Publish Version
= 端到端可见延迟
监控只观察 Doris 写入耗时,无法解释 Checkpoint 排队导致的延迟。
七、Flink CDC 整库同步与 Doris 更新语义

数据库 CDC 需要把操作日志转换为 Doris 当前状态。官方数据更新文档给出了一组协同机制:Unique Key UPSERT 处理 INSERT 和 UPDATE,__DORIS_DELETE_SIGN__ 处理 DELETE,Sequence Column 使用 Binlog 时间戳、LSN 或其他单调版本解决乱序。
1. INSERT 与 UPDATE
目标表使用 Unique Key。相同主键再次写入时形成新版本,Merge-on-Write 在写入阶段维护最新状态和 Delete Bitmap。源库更新主键时,Debezium/Flink CDC 通常会生成旧主键 DELETE 和新主键 INSERT,两条事件需要完整传递。
2. DELETE
删除事件不能被当成“没有数据”。Connector 需要开启删除同步,或显式生成 __DORIS_DELETE_SIGN__。删除生效后,查询会过滤对应旧版本;物理空间由后续 Compaction 回收。
3. 乱序与 Sequence Column
培训材料中的订单案例很典型:10:01 的 SHIPPED 先到 Doris,10:00 的 PAID 因网络延迟后到。如果只按照到达顺序覆盖,最终状态会回退到 PAID。Sequence Column 使用源端提交时间或日志位点比较版本,只接受更大的 Sequence。
Sequence 值需要满足:
- 同一业务主键上可以比较新旧;
- 重试不会改变事件原始 Sequence;
- 跨时区和时间精度明确;
- 时钟回拨场景有预案;
- 多数据流共同更新一行时,评估 Sequence Mapping 或字段级更新方案。
Exactly-Once 控制重复提交,Sequence Column 控制版本新旧。两者解决的问题不同,生产 CDC 表通常需要同时设计。
4. 部分列更新
源事件只包含变化字段时,可以使用 Unique Key 部分列更新。部分更新会增加历史值读取、补齐完整行和随机 I/O,宽表高频更新需要在真实 SSD、目标并发和目标字段分布下压测。灵活部分列更新、固定列部分更新和整行 UPSERT 应分开选择,避免把所有 CDC 流统一套用一个模板。
八、整库同步与 Schema 演进

Flink Doris Connector 集成 Flink CDC,可以完成 MySQL、Oracle、PostgreSQL 等数据库的整库或多表同步,并支持自动建表和 DDL 同步。它适合表数量多、字段变化频繁、团队已经能够运维 Flink 的场景。
1. 全量到增量的切换
典型流程是:
- Flink CDC 读取源表 Snapshot;
- 同时记录与 Snapshot 一致的日志切换点;
- Connector 创建或检查 Doris 目标表;
- 全量数据按主键写入;
- Snapshot 完成后从记录的 Binlog/WAL 位点继续;
- Checkpoint 持续保存 Source State 和 Sink State。
全量阶段的读取压力会影响源库,通常需要设置分片键、并发、限速和运行窗口。无主键表、超大单表、热点主键和大字段表需要单独评估。
2. 自动建表带来的收益与风险
自动建表减少人工 DDL,整库同步可以更快启动。它也会把上游类型映射、主键判断、默认副本数和分桶策略带入下游。自动生成的表结构很难同时满足所有查询模式。生产项目通常采用两级策略:
- 普通镜像表允许自动创建,再通过表模板控制副本和基础属性;
- 核心大表预先创建,明确数据模型、分区、分桶、排序键、索引和 Sequence Column。
3. DDL 同步需要真实回归集
“支持 Schema Evolution”不能等同于支持任意 DDL。上游常见变化包括:
- 新增可空列;
- 修改列默认值;
- 扩大字符串长度;
- 数值类型扩宽;
- 改名、删除列;
- 修改主键;
- 修改列类型;
- 表重建与交换。
不同 Connector 版本、Flink CDC 版本和 Doris 版本支持范围会变化。上线前应从生产审计日志中抽取真实 DDL,建立自动回归测试。未支持的 DDL 需要进入暂停、告警、人工建表、回填和恢复流程。
4. Connector 版本与升级治理
当前 Flink Doris Connector 26.1.1 支持多组 Flink 版本。生产升级需要同时校验:
- Connector Jar 与 Flink 版本;
- Java 运行时;
- Flink State Serializer 兼容性;
- Checkpoint/Savepoint 恢复;
- Stream Load 参数和 Doris 版本;
- CDC Source Connector 与数据库版本;
- DDL 类型映射;
- 删除、部分更新和主键变更;
- TLS、认证和网络策略。
升级演练应从 Savepoint 启动新 Job,验证原 Job 停止点、新 Job 起始点和 Doris 最终结果。
九、Doris Streaming Job:把持续同步能力放入 Doris

Doris 4.x 的 Streaming Job 提供内置持续导入能力,可以持续读取 MySQL、PostgreSQL 和 S3。用户通过 CREATE JOB ... ON STREAMING 创建任务,不需要单独维护一套外部 Flink 集群。当前数据库 CDC 主要有两种模式。
1. SQL Mapping Sync
SQL Mapping 通过 Job + CDC Stream TVF 实现,目标表需要提前创建。用户使用 INSERT INTO ... SELECT 明确字段映射、过滤和类型转换。
CREATE JOB day10_mysql_order_sync
ON STREAMING
DO
INSERT INTO day10_lab.order_state
(order_id, user_id, status, amount, update_time, source_commit_ts)
SELECT
id,
buyer_id,
upper(order_status),
cast(pay_amount AS DECIMAL(18,2)),
update_time,
row_version
FROM cdc_stream(
"type" = "mysql",
"jdbc_url" = "jdbc:mysql://mysql-host:3306",
"driver_url" = "mysql-connector-java-8.0.25.jar",
"driver_class" = "com.mysql.cj.jdbc.Driver",
"user" = "cdc_user",
"password" = "${MYSQL_PASSWORD}",
"database" = "biz_db",
"table" = "orders",
"offset" = "initial"
)
WHERE pay_amount >= 0;
SQL Mapping 支持 Exactly-Once,适合单表、目标表已经精细设计、需要字段裁剪和转换的链路。目前同步源表需要主键,MySQL 需要开启 ROW 格式 Binlog,PostgreSQL 需要开启 Logical Replication。
2. Auto Table Creation Sync
自动建表模式可以同步一组表或整个数据库:
CREATE JOB day10_mysql_db_sync
ON STREAMING
FROM MYSQL (
"jdbc_url" = "jdbc:mysql://mysql-host:3306",
"driver_url" = "mysql-connector-java-8.0.25.jar",
"driver_class" = "com.mysql.cj.jdbc.Driver",
"user" = "cdc_user",
"password" = "${MYSQL_PASSWORD}",
"database" = "biz_db",
"include_tables" = "users,orders,order_items",
"offset" = "initial"
)
TO DATABASE day10_mirror (
"table.create.properties.replication_num" = "3"
);
offset=initial 表示先执行全量初始化,再切换增量;latest 表示从最新位置开始。当前自动建表同步保证 At-Least-Once,只支持主键表,不提供列映射、过滤和转换。目标表已存在时会跳过创建,因此核心表可以提前按生产标准建好。
3. S3 Continuous Load
Streaming Job 还可以持续监控 S3 路径,发现新文件后执行 INSERT INTO ... SELECT FROM S3(...)。这种模式适合上游周期性产出文件、又希望避免外部调度器轮询的场景。任务状态可以通过 jobs("type"="insert") 表函数查看,CurrentOffset 可反映已处理文件进度。
4. 何时选择 Streaming Job
它的优势是组件更少、SQL 表达清晰、任务与 Doris 权限治理统一。生产决策仍需评估:
- 当前目标版本是否完整支持所需源端;
- 源端数据库规模和日志吞吐;
- 全量初始化并行度与源库压力;
- 复杂转换、Join 和窗口计算需求;
- 团队是否需要 Flink 的成熟生态与状态运维工具;
- 任务数量、并发和故障隔离;
- 驱动 Jar、密码和密钥管理;
- 版本升级与回滚路径。
十、把各种“一致性”放在同一张表里
流式方案评审中,经常出现所有产品都写着“可靠”,却没有说明具体保证。下面的矩阵可以直接放入架构评审文档。
| 方案 | 状态位置 | 默认/主要语义 | 重复处理 | 乱序处理 | 删除处理 |
|---|---|---|---|---|---|
| Routine Load | Doris FE 保存 Kafka Offset | Kafka → Doris Exactly-Once | Task 事务与提交校验 | Sequence Column | WITH DELETE / Delete Sign |
| Kafka Connector | Kafka Connect Offset | 取决于 Connector 配置 | Kafka Connect + 目标模型 | Sequence Column | 事件映射 |
| Flink Connector Streaming | Flink Checkpoint | Exactly-Once | Stream Load 2PC | Sequence Column | sink.enable-delete |
| Flink Connector Batch | Connector 合批状态 | At-Least-Once | Unique Key 幂等 | Sequence Column | Delete Sign |
| Flink CDC 整库同步 | Flink CDC State | Checkpoint 驱动 | 2PC + Unique Key | Source Position / Sequence | CDC RowKind |
| Streaming Job SQL Mapping | Doris Job State | Exactly-Once | Job + CDC Stream TVF | Sequence/源位点 | CDC 事件 |
| Streaming Job 自动建表 | Doris Job State | At-Least-Once | Unique Key | 源位点与目标表设计 | CDC 事件 |
需要特别强调三点:
第一,Exactly-Once 无法修复源端已经重复产生的业务事件。业务事件 ID 和源系统约束仍然重要。
第二,Unique Key 去重只对相同主键有效。主键选择错误会把两条不同业务记录合并,也可能让同一业务对象产生多个 Key。
第三,乱序由 Sequence 机制处理。没有顺序列时,最后到达的事件会成为最后状态,它与源端真实提交顺序可能不同。
十一、监控、告警与故障排查

一条实时链路至少需要五层监控。
1. 状态层
- Routine Load Job State;
- Flink Job State 与 Restart 次数;
- Streaming Job Status;
- 暂停原因、异常消息和最近成功时间。
2. 进度层
- Kafka 每个 Partition 的当前 Offset 和 End Offset;
- Consumer Lag 与 Lag 增长速度;
- Flink Checkpoint Age、Duration 和失败率;
- CDC Binlog/LSN 延迟;
- Streaming Job CurrentOffset;
- 业务事件时间到 Doris 可见时间的差值。
3. 吞吐与资源层
- Rows/s、Bytes/s、Task 数;
- FE 事务和元数据压力;
- BE CPU、内存、磁盘 I/O、网络;
- Load MemTable 和 Flush;
- Compaction Score、Rowset 和 Segment 数;
- 目标表 Tablet 分布与热点。
4. 数据质量层
- 过滤行数和过滤比例;
- 类型转换错误、JSON 解析错误;
- NULL 约束和长度超限;
- 源表与目标表字段差异;
- 删除事件和顺序列缺失;
- 异常样本的隔离与重放状态。
5. 业务结果层
- 源端与 Doris 行数;
- 主键集合差异;
- 金额、数量和状态分布;
- 最大更新时间;
- 删除数量;
- 指定业务批次的逐主键抽样。
常见故障及处理路径
Lag 持续上升。 先看任务是否运行,再看实际并发、Kafka Partition、BE 资源、批次阈值、过滤率、索引写入成本和 Compaction。必要时从源端生产速率与下游最大稳定消费速率计算缺口。
Routine Load 进入 PAUSED。 查看 ReasonOfStateChanged 和 OtherMsg。数据质量导致的暂停通常不会自动恢复,应修正消息或调整质量规则;认证、网络和 Kafka 暂时不可用可以在故障解除后恢复。
Offset out of range。 Kafka 已清理待消费消息。需要确认可用最早 Offset,评估从归档、源库或对象存储补数,然后重设消费起点。直接跳到最新 Offset 会留下数据缺口,必须形成业务审批和缺口记录。
Flink 重启后 Label 冲突。 检查是否从最新 Checkpoint/Savepoint 恢复,检查 sink.label-prefix 是否被其他 Job 使用。改变 Prefix 会形成新的事务命名空间,需要配合对账。
状态出现回退。 检查 Sequence Column、时区、精度和源端日志位点。相同主键的 Sequence 值需要单调可比较。
Schema 变更导致失败。 暂停受影响表或 Job,确认上游 DDL、目标表结构、Connector 支持矩阵和类型映射。修复后从保存的日志位点恢复,并执行新旧 Schema 的数据对账。
恢复完成需要满足四个条件:任务状态恢复、Lag 持续收敛、端到端 P99 回到 SLA、数据对账通过。
十二、方案选择方法

可以使用下面这条路径完成评审:
第一步:确定源端和消息格式
Kafka CSV/JSON 可以直接进入 Routine Load;Avro/Protobuf/Debezium 评估 Kafka Connector;数据库日志进入 Flink CDC 或 Streaming Job。
第二步:确定处理复杂度
只做字段映射、过滤和派生列,优先选择更短链路。存在窗口、Join、维表、复杂状态和多流编排时,使用 Flink。
第三步:确定一致性合同
写明进度保存位置、提交边界、故障恢复点、重复处理、乱序处理和删除语义。不要只写一个“Exactly-Once”标签。
第四步:确定表结构
CDC 当前状态表通常采用 Unique Key MoW。主键、Sequence Column、分区、分桶、索引和部分更新策略需要一起评审。原始事件可以另存 Duplicate Key 明细表,用于审计和重放。
第五步:确定容量和 SLA
至少提供:
- 源端峰值 Rows/s 与 Bytes/s;
- 日增量和历史初始化数据量;
- Kafka Partition 或 CDC Snapshot 并发;
- 目标可见延迟 P50/P95/P99;
- 故障允许恢复时间;
- Kafka/Binlog/WAL 保留时间;
- Doris BE 数量、磁盘和网络;
- Flink State 大小和 Checkpoint 时长;
- 目标表 Tablet 数、索引和 Compaction 预算。
第六步:完成故障演练
演练 Kafka Broker 不可用、BE 重启、FE 切主、Flink TaskManager 失败、Checkpoint 失败、网络断开、坏数据、Schema 加列、Offset 过期和源端日志不足。每次演练记录恢复点、恢复耗时、重复数、缺失数和人工步骤。
第七步:定义 POC 验收指标
实时同步 POC 不能只看“任务是否跑起来”。建议至少运行 24–72 小时,并记录下列指标:
| 维度 | 验收内容 |
|---|---|
| 正确性 | 全量行数、主键集合、聚合值、删除数、状态分布 |
| 延迟 | 事件时间到 Doris 可见的 P50、P95、P99 |
| 吞吐 | 平均与峰值 Rows/s、Bytes/s,Lag 是否收敛 |
| 恢复 | Kafka、Flink、FE、BE 故障后的恢复点和恢复时间 |
| 重复 | 重放、超时和 Checkpoint 失败后是否出现异常重复 |
| 乱序 | 晚到旧事件是否被 Sequence 正确拒绝 |
| 变更 | 新增列、类型变化和表创建是否符合支持矩阵 |
| 资源 | FE 事务、BE 内存、磁盘、网络与 Compaction 水位 |
压测数据要尽量接近真实分布。平均消息大小、热点主键比例、单行最大尺寸、更新占比、删除占比、宽表列数和索引数量都会改变结果。POC 报告应保留源端生产速率、Kafka Partition、Flink 并发、Checkpoint、Doris 表 DDL、集群规格和任务参数,保证结论可以复现。
十三、Day 10 实战任务
完整包已经提供 SQL、Flink SQL、样例数据和运维模板。建议按下面顺序完成。
实验一:Routine Load 基础链路
- 创建
order_event_detailDuplicate Key 表和order_stateUnique Key 表; - 创建 Kafka CSV 与 JSON Topic;
- 写入正常订单事件;
- 创建两个 Routine Load Job;
- 使用
SHOW ROUTINE LOAD记录 State、Progress 和 Statistic; - 验证数据可见延迟和目标行数。
实验二:暂停、恢复与 Offset
- 手工 Pause Job;
- 在暂停期间继续生产 Kafka 消息;
- 记录 End Offset 和 Doris 当前 Offset;
- Resume Job;
- 观察 Lag 收敛;
- 验证暂停期间消息全部进入 Doris。
实验三:数据质量
- 发送金额格式错误、缺少主键和非法时间消息;
- 观察 Job 是否暂停;
- 检查错误信息和异常样本;
- 修复消息或质量规则;
- 恢复 Job;
- 记录过滤比例和补数过程。
实验四:乱序更新
依次构造:
T1: order_id=1001, status=PAID, sequence=100
T2: order_id=1001, status=SHIPPED, sequence=200
让 T2 先到、T1 后到,验证最终状态保持 SHIPPED。再创建一个没有 Sequence Column 的对照表,观察结果差异。
实验五:Flink Checkpoint 2PC
- 部署匹配版本的 Flink Doris Connector;
- 创建 Datagen 或 MySQL CDC Source;
- 配置 10 秒 Checkpoint;
- 开启 2PC 和唯一 Label Prefix;
- 在写入期间杀死 TaskManager;
- 从 Checkpoint 恢复;
- 对比源端主键集合与 Doris 状态。
实验六:Streaming Job 模板验证
- 预建一个 Unique Key 目标表;
- 创建 MySQL SQL Mapping Job;
- 验证全量初始化;
- 执行 INSERT、UPDATE、DELETE;
- 查看
jobs("type"="insert"); - 记录 CurrentOffset、状态和对账结果。
最终提交《Doris 实时同步方案与验收单》,内容至少包括:源端、数据量、延迟目标、语义范围、表结构、配置、监控、异常处理、回滚和验证证据。
十四、生产最佳实践清单
设计阶段
- 明确业务主键、事件 ID 和 Sequence 来源;
- 区分事件明细表与当前状态表;
- 把 Exactly-Once 的起点和终点写清楚;
- 确定 Kafka/Binlog/WAL 保留时间;
- 为全量初始化设置源库限流;
- 建立 DDL 支持矩阵和人工处置流程;
- 核心大表优先预建,避免自动建表产生不合适的分桶和排序;
- 敏感密码和 AK/SK 进入密钥管理系统。
上线阶段
- 使用独立写入账号和最小权限;
- Connector、Flink、Java 和 Doris 版本组合完成回归;
- Label Prefix 全局唯一;
- Checkpoint 存储高可用;
- 监控 Lag、失败率、过滤率、Checkpoint 和 Compaction;
- 进行 Pause/Resume、重启、网络和 Schema 故障演练;
- 完成全量和增量双重对账;
- 设定清晰的上线、暂停和回滚门槛。
安全与权限
流式任务通常长期保存数据库密码、Kafka 认证信息和 Doris 写入凭据。生产环境应使用独立服务账号,权限限制在指定数据库和表;密码放入 Secret、Vault 或集群密钥管理系统;日志和 SHOW 输出需要避免泄露明文。MySQL CDC 账号只授予读取 Snapshot 与 Binlog 所需权限,PostgreSQL 账号只授予逻辑复制和目标表读取权限。跨机房链路应启用 TLS,并在证书轮换前完成 Connector 兼容测试。
账号变更、证书轮换和密码更新都属于运行变更。操作前记录 Offset 或 Savepoint,更新后验证任务状态、Lag 和对账结果。
运行阶段
- 定期检查长期 Lag 趋势和峰值消费余量;
- 监控 Kafka Retention 与 CDC 日志保留;
- 把异常样本、任务状态和对账结果写入审计表;
- 大版本升级从 Savepoint 或明确 Offset 演练;
- Schema 变更进入发布流程,禁止无通知直接改核心源表;
- 对热点主键、超大行和高频部分更新单独治理;
- 定期验证 Runbook,确保值班人员可以执行恢复。
十五、Knowledge Check
1. Routine Load 的长期 Job 和短期 Task 分别负责什么?
答案要点: Job 管理生命周期、状态、配置、进度和 Task 拆分;Task 负责消费一组 Kafka Partition、执行一批数据写入,并以独立事务提交。
2. Routine Load 的实际 Task 数由哪些因素决定?
答案要点: Kafka Topic Partition 数、desired_concurrent_number、FE 配置的 Routine Load Task 并发上限三者取最小值。
3. Routine Load 怎样实现 Exactly-Once?
答案要点: 事务提交时持久化消费进度,并在提交前检查 Task 是否仍有效,避免旧 Task 或重复 Task 形成第二次有效提交。
4. PAUSED 与 STOPPED 有什么区别?
答案要点: PAUSED 仍可 Resume,部分异常可以自动恢复;STOPPED 是终态,无法恢复。
5. Flink Connector 的流式写入为何能实现 Exactly-Once?
答案要点: 以 Flink Checkpoint 为边界,使用 Stream Load 2PC Prepare/Commit,把 Flink State 成功与 Doris 版本发布对齐。
6. sink.label-prefix 有什么要求?
答案要点: 2PC 场景下需要全局唯一且稳定,用于构造事务 Label。恢复时应从正确 Checkpoint/Savepoint 启动。
7. CDC 删除事件怎样进入 Doris?
答案要点: Flink Connector 可以根据 RowKind 写入 Delete Sign,也可以由应用显式生成 __DORIS_DELETE_SIGN__;目标表通常使用 Unique Key。
8. Exactly-Once 已经开启,为什么仍需要 Sequence Column?
答案要点: Exactly-Once 控制同一批数据的重复有效提交;Sequence Column 比较同一主键事件的新旧,防止旧事件晚到后覆盖新状态。
9. Streaming Job 的两种数据库同步模式有什么差异?
答案要点: SQL Mapping 需要预建目标表,支持映射、过滤、转换和 Exactly-Once;自动建表模式适合整库镜像,当前保证 At-Least-Once,不支持映射和转换。
10. 流式任务恢复成功需要哪些证据?
答案要点: 任务状态正常、Lag 持续收敛、端到端延迟回到 SLA、源端与 Doris 对账通过,并记录恢复点与异常处理过程。
本日小结
Day 10 把“持续写入”拆成了几套可以验证的机制:Routine Load 在 Doris 内部保存 Kafka Offset,并用长期 Job 管理不断生成的微批 Task;Flink Doris Connector 把 Checkpoint 与 Stream Load 2PC 对齐;Flink CDC 负责 Snapshot、日志位点和 Schema 变化;Unique Key、Delete Sign 和 Sequence Column共同维护下游当前状态;Streaming Job 提供更短的内置持续同步路径。
一条实时链路的质量可以用五句话检查:进度有持久化位置,提交有明确边界,重复有去重机制,乱序有版本规则,故障有恢复与对账证据。
Day 11 将进入 Doris 索引体系,从 Prefix Index、ZoneMap、Bloom Filter、Bitmap、Inverted Index 和 NGram Bloom Filter 出发,理解 Doris 怎样通过 Data Skipping 减少实际扫描数据。
官方资料
- Apache Doris 4.x|Routine Load
- Apache Doris 4.x|Routine Load Principles and Best Practices
- Apache Doris 4.x|Importing Data from Kafka
- Apache Doris|Flink Doris Connector
- Apache Doris 4.x|Data Update Overview
- Apache Doris 4.x|CDC_STREAM TVF
- Apache Doris 4.x|CREATE STREAMING JOB
- Apache Doris 4.x|MySQL CDC with Auto Table Creation
- Apache Doris 4.x|PostgreSQL CDC with SQL Mapping
- Apache Doris Download