Doris 实战与架构

分析型数据库 / 源码阅读 / LESSON 20

实时写入、更新与删除——CDC 场景的正确打开方式

围绕 CDC、Label 幂等、版本堆积、Unique Key、删除语义和 Routine Load 建立实时写入方法论。

阅读时间
50 分钟
学习路径
Doris 实战与架构
内容来源
Doris 深度笔记

预计阅读时间: 50 分钟 前置阅读: doc-02, doc-07, doc-10, doc-17 下一次阅读: doc-21(湖仓与联邦查询)


1. 实时不是只看导入延迟

一个 Doris 实时链路至少包含五段延迟:

源系统产生数据
  → 采集/CDC/消息队列
    → Doris Load 事务
      → Publish Version 可见
        → 查询命中最新版本

专家判断实时性时, 不只看 Stream Load 返回多快, 还要看:

  • 是否有小批量高频写入造成版本堆积。
  • 是否有 Compaction 跟不上。
  • Unique Key 更新是否放大查询成本。
  • Kafka/Routine Load 是否出现消费延迟。
  • 下游查询是否要求强一致读最新。

参考:


2. 导入方式的生产边界

方式适合边界
Stream Load应用直接写入、小批文件、准实时客户端要控制批大小和 Label
Routine LoadKafka 持续消费要管消费延迟、错误数据和分区变化
Broker LoadHDFS/S3 大批文件分钟级任务, 不适合低延迟
INSERT INTO小批 SQL 写入不适合大批量 VALUES
Group Commit高频小 INSERT 合并需要接受同步/异步模式差异
Flink/CDC源库变更同步主键、乱序、删除语义必须确认

经验:

  • 大批量不要用一条超大 Load, 失败重试成本太高。
  • 高频小批不要直接把每条消息打到 Doris。
  • 能在客户端或流计算层合并, 就先合并。
  • 使用 Label 保证幂等, 避免重试造成重复。

3. Label 与幂等

Stream Load 的 Label 是导入幂等的核心。

curl --location-trusted -u user:password \
  -H "label: order_20260727_0001" \
  -H "format: json" \
  -H "read_json_by_line: true" \
  -T orders.json \
  http://fe:8030/api/db/order_detail/_stream_load

同一个 Label 成功后再次提交, Doris 会识别重复。生产策略:

source + partition + offset-range + attempt

不要用随机 UUID 当唯一 Label, 否则重试无法幂等。


4. 版本堆积与 Compaction 压力

每次成功写入都会产生新版本。频率过高的小批量导入会带来:

  • Rowset 数量增加。
  • 查询需要合并更多版本。
  • Compaction 压力上升。
  • Tablet 元数据和调度开销增加。

检查:

SHOW TABLETS FROM table_name;

如果 Version Count 很高:

  1. 增大每批写入数据量。
  2. 降低提交频率。
  3. 使用 Group Commit。
  4. 检查 Compaction 线程和磁盘压力。
  5. 评估是否拆分热点表或调整 Bucket。

5. Unique Key 与 Merge-on-Write

Unique Key 适合 CDC, 但要先确认三件事:

  1. 主键是否稳定且足够均匀。
  2. 更新频率是否可控。
  3. 查询是否主要读最新状态。

建表示例:

CREATE TABLE order_current
(
    order_id BIGINT NOT NULL,
    tenant_id BIGINT NOT NULL,
    status TINYINT,
    amount DECIMAL(18, 2),
    updated_at DATETIME,
    delete_sign TINYINT DEFAULT "0"
)
UNIQUE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 32
PROPERTIES (
    "enable_unique_key_merge_on_write" = "true",
    "replication_num" = "3"
);

Merge-on-Write 的好处是查询读到的就是合并后结果, 读性能更稳定。代价是写入阶段承担更多合并工作。

反例:

  • 用 Unique Key 存行为日志。
  • 主键包含时间戳, 导致根本不会更新。
  • 高频更新大宽表, 但只查询少数字段。

6. 删除语义

Doris 删除有多种形态:

方式适合风险
分区删除按时间 TTL 删除最干净, 但粒度较粗
DELETE 条件删除少量条件删除可能产生删除标记和 Compaction 压力
Unique Key 删除标记CDC 删除要统一上游 delete 语义
TRUNCATE 分区/表清空重跑不可恢复, 需要审批

如果业务可以按时间生命周期删除, 优先设计 Partition TTL, 不要每天执行大量条件 DELETE。


7. Routine Load 与 Kafka

Routine Load 的关键不是创建语法, 而是消费治理:

CREATE ROUTINE LOAD db.routine_order
ON order_detail
COLUMNS(...)
PROPERTIES
(
    "desired_concurrent_number" = "8",
    "max_batch_interval" = "10",
    "max_batch_rows" = "500000",
    "max_batch_size" = "104857600"
)
FROM KAFKA
(
    "kafka_broker_list" = "broker1:9092,broker2:9092",
    "kafka_topic" = "order_topic"
);

关注:

  • Kafka 分区数与 Doris 并发是否匹配。
  • 错误数据比例是否触发暂停。
  • 消费 Offset 是否符合重放预期。
  • BE 不可用时任务是否堆积。
  • Topic 增加分区后 Routine Load 是否识别。

8. CDC 到 Doris 的专家清单

1. 源库主键是否完整
2. 更新乱序如何处理
3. Delete 事件如何表达
4. DDL 变更是否同步
5. Label 或 Offset 如何保证幂等
6. 小批量写入如何合并
7. Unique Key 是否会造成热点
8. 历史回补和实时增量是否使用同一张表
9. 数据校验如何做
10. 失败重放从哪里开始

数据校验示例:

-- 行数校验
SELECT count(*) FROM order_current;

-- 主键重复校验, Unique Key 正常情况下不应有重复可见行
SELECT order_id, count(*)
FROM order_current
GROUP BY order_id
HAVING count(*) > 1;

-- 延迟校验
SELECT max(updated_at) FROM order_current;

9. 实时链路的 SLO

不要只写"秒级实时", 要写可测 SLO:

指标示例
P95 可见延迟30 秒内
最大可接受延迟5 分钟
重试后重复率0
错误数据隔离进入错误表或错误文件
回补窗口支持最近 7 天重放
查询新鲜度dashboard 显示数据时间

一句话总结:

Doris 实时写入专家不是追求每条数据立刻入库, 而是在幂等、批量、版本、Compaction、更新语义和查询新鲜度之间做可验证的工程平衡。