Doris architecture

ANALYTICAL DATABASE / SOURCE READING / LESSON 20

Real-time ingestion, updates, and deletes

Build a production method for CDC, idempotent labels, version buildup, Unique Key updates, delete semantics, and Routine Load.

Reading
50 min
Track
Doris architecture
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 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、更新语义和查询新鲜度之间做可验证的工程平衡。