第 10 关 · ★★★

流式同步与 CDC

Routine Load、Flink Connector、Streaming Job:把上游变更持续同步进 Doris。

已点亮 · 最佳 分

流式同步与 CDC:Routine Load、Flink Connector 与 Streaming Job

从「每秒能写多少行」到「故障后能证明什么」——设计一条能进生产的实时链路

Routine LoadExactly-OnceFlink 2PCCDC 语义Streaming Job
1 / 28

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

Day 10 学习总览
Day 10 学习总览

Day 9 解决了批量数据怎样进入 Doris:客户端可以通过 Stream Load 推送文件,Doris 可以通过 TVF、Catalog 或历史上的 Broker Load 读取远端数据,内部表之间也可以使用 INSERT INTO ... SELECT 完成加工。批量任务拥有清晰的开始和结束,一次提交对应一个相对明确的数据集合。

流式同步把问题延伸到了时间轴上。Kafka 会持续产生消息,MySQL Binlog 和 PostgreSQL WAL 会不断记录变化,业务表可能新增字段,网络会抖动,任务会重启,数据会重放,旧事件也可能晚于新事件抵达。系统需要长期保存读取进度,并在每次恢复时找到正确的继续位置。这里的“实时”同时包含四层含义:数据持续进入、状态持续推进、故障后能够恢复、查询结果持续逼近源端真实状态。

因此,流式链路的设计重点不止是“每秒能写多少行”。一套可以进入生产的方案,还要回答下面这些问题:

  • Kafka Offset、Flink Checkpoint 或数据库日志位点由谁维护;
  • 一批数据写入成功时,读取进度是否同步提交;
  • 任务在提交前宕机、提交后丢失响应时会发生什么;
  • INSERTUPDATEDELETE 怎样映射到 Doris 表模型;
  • 旧事件晚到时,怎样防止它覆盖新状态;
  • 源表加列、改类型或改主键时,链路怎样处理;
  • 任务恢复后,怎样证明无缺失、无异常重复、状态无回退;
  • 延迟升高时,问题位于源端、流处理、Doris 写入,还是后台 Compaction。

本篇将围绕 Routine Load、Doris Kafka Connector、Flink Doris Connector、Flink CDC 和 Doris Streaming Job 建立一张完整地图。重点放在原理、语义、选型和生产验收。完成本篇后,学习者应当能够独立设计一条 Kafka 或数据库 CDC 到 Doris 的实时同步链路,并写出相应的运行手册与验收标准。

本日学习目标

完成本篇后,你应当能够:

  1. 根据数据源、格式、处理复杂度和一致性要求选择实时同步方式;
  2. 解释 Routine Load 的 Job、Task、Kafka Partition 和 Doris 事务之间的关系;
  3. 说明 Routine Load 如何通过消费进度持久化和提交校验实现 Exactly-Once;
  4. 正确设置起始 Offset、任务并发和微批阈值;
  5. 使用 SHOW ROUTINE LOAD 判断任务状态、消费进度和暂停原因;
  6. 解释 Flink Checkpoint 与 Stream Load 2PC 怎样共同控制数据可见性;
  7. 使用 Unique Key、Delete Sign 和 Sequence Column 承接 CDC 事件;
  8. 理解 Flink CDC 整库同步、自动建表和 Schema 演进的边界;
  9. 区分 Streaming Job 的 SQL Mapping Sync 与 Auto Table Creation Sync;
  10. 建立覆盖 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 架构
Routine Load 架构

你提供的培训材料用一张图概括了 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:

  1. 负责一组 Kafka Partition;
  2. 从已记录 Offset 开始读取;
  3. 在达到时间、行数或数据量阈值时结束当前微批;
  4. 把数据写入 Doris;
  5. 以独立事务提交;
  6. 提交成功后更新消费进度;
  7. 由 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 消息包含 beforeafterop 和源端位点等结构,进入 Doris 前要明确主键、事件操作类型、字段映射与顺序规则。使用 Kafka Connector 并不会自动消除业务乱序,Sequence Column 和源端版本仍要纳入表设计。

选型时还要考虑运维归属。Routine Load 的运行状态和消费进度在 Doris 中查看;Kafka Connector 的 Worker、Connector、Task 和 Offset 需要在 Kafka Connect 平台中监控。企业已经有成熟 Kafka Connect 平台时,后者可以复用现有治理;团队希望减少中间组件时,Routine Load 更轻。


三、Routine Load 怎样实现 Exactly-Once

Routine Load Exactly-Once
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_intervalmax_batch_rowsmax_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 状态机
Routine Load 状态机

Routine Load 的主要状态包括:

状态 含义 后续动作
NEED_SCHEDULE 等待初次调度或重新调度 调度成功后进入 RUNNING
RUNNING 正常持续消费 可能暂停或停止
PAUSED 任务暂停,未终止 手工 Resume 或满足条件后自动恢复
STOPPED 手工永久停止 无法恢复
CANCELLED 表、库被删除或出现不可恢复错误 无法恢复

Kafka 暂时不可用、网络抖动等异常通常能够触发自动恢复。手工执行 PAUSE、数据质量问题和不可恢复错误一般需要人工处理。恢复动作前应先检查 ReasonOfStateChanged,直接反复 Resume 可能让任务在错误条件未消除时持续抖动。

Offset 与微批调节
Offset 与微批调节

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 升高时的排查顺序

  1. 确认 Job 处于 RUNNING
  2. 对比 Kafka Partition 数和实际 Task 数;
  3. 检查 desired_concurrent_number 与 FE 上限;
  4. 检查 BE CPU、内存、磁盘和网络;
  5. 检查过滤率、解析错误和单行超大消息;
  6. 检查目标表索引、Unique Key 更新和 Delete Bitmap 成本;
  7. 检查 Rowset、Segment 与 Compaction Score;
  8. 检查 Kafka Broker 吞吐和跨机房网络。

只提高并发可能把瓶颈从消费端移动到 BE Flush 或 Compaction。调优需要同时观察端到端延迟和存储健康度。


六、Flink Doris Connector:把 Flink State 与 Doris 版本对齐

Flink Checkpoint 与 Stream Load 2PC
Flink Checkpoint 与 Stream Load 2PC

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 一致性合同
CDC 一致性合同

数据库 CDC 需要把操作日志转换为 Doris 当前状态。官方数据更新文档给出了一组协同机制:Unique Key UPSERT 处理 INSERTUPDATE__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 演进

整库同步与 Schema 演进
整库同步与 Schema 演进

Flink Doris Connector 集成 Flink CDC,可以完成 MySQL、Oracle、PostgreSQL 等数据库的整库或多表同步,并支持自动建表和 DDL 同步。它适合表数量多、字段变化频繁、团队已经能够运维 Flink 的场景。

1. 全量到增量的切换

典型流程是:

  1. Flink CDC 读取源表 Snapshot;
  2. 同时记录与 Snapshot 一致的日志切换点;
  3. Connector 创建或检查 Doris 目标表;
  4. 全量数据按主键写入;
  5. Snapshot 完成后从记录的 Binlog/WAL 位点继续;
  6. 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

Streaming Job 两种模式
Streaming Job 两种模式

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。 查看 ReasonOfStateChangedOtherMsg。数据质量导致的暂停通常不会自动恢复,应修正消息或调整质量规则;认证、网络和 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 基础链路

  1. 创建 order_event_detail Duplicate Key 表和 order_state Unique Key 表;
  2. 创建 Kafka CSV 与 JSON Topic;
  3. 写入正常订单事件;
  4. 创建两个 Routine Load Job;
  5. 使用 SHOW ROUTINE LOAD 记录 State、Progress 和 Statistic;
  6. 验证数据可见延迟和目标行数。

实验二:暂停、恢复与 Offset

  1. 手工 Pause Job;
  2. 在暂停期间继续生产 Kafka 消息;
  3. 记录 End Offset 和 Doris 当前 Offset;
  4. Resume Job;
  5. 观察 Lag 收敛;
  6. 验证暂停期间消息全部进入 Doris。

实验三:数据质量

  1. 发送金额格式错误、缺少主键和非法时间消息;
  2. 观察 Job 是否暂停;
  3. 检查错误信息和异常样本;
  4. 修复消息或质量规则;
  5. 恢复 Job;
  6. 记录过滤比例和补数过程。

实验四:乱序更新

依次构造:

T1: order_id=1001, status=PAID,    sequence=100
T2: order_id=1001, status=SHIPPED, sequence=200

让 T2 先到、T1 后到,验证最终状态保持 SHIPPED。再创建一个没有 Sequence Column 的对照表,观察结果差异。

实验五:Flink Checkpoint 2PC

  1. 部署匹配版本的 Flink Doris Connector;
  2. 创建 Datagen 或 MySQL CDC Source;
  3. 配置 10 秒 Checkpoint;
  4. 开启 2PC 和唯一 Label Prefix;
  5. 在写入期间杀死 TaskManager;
  6. 从 Checkpoint 恢复;
  7. 对比源端主键集合与 Doris 状态。

实验六:Streaming Job 模板验证

  1. 预建一个 Unique Key 目标表;
  2. 创建 MySQL SQL Mapping Job;
  3. 验证全量初始化;
  4. 执行 INSERT、UPDATE、DELETE;
  5. 查看 jobs("type"="insert")
  6. 记录 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. PAUSEDSTOPPED 有什么区别?

答案要点: 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 减少实际扫描数据。


官方资料

  1. Apache Doris 4.x|Routine Load
  2. Apache Doris 4.x|Routine Load Principles and Best Practices
  3. Apache Doris 4.x|Importing Data from Kafka
  4. Apache Doris|Flink Doris Connector
  5. Apache Doris 4.x|Data Update Overview
  6. Apache Doris 4.x|CDC_STREAM TVF
  7. Apache Doris 4.x|CREATE STREAMING JOB
  8. Apache Doris 4.x|MySQL CDC with Auto Table Creation
  9. Apache Doris 4.x|PostgreSQL CDC with SQL Mapping
  10. Apache Doris Download