Doris 的数据从哪来:Kafka/RocketMQ 导入的几种姿势与取舍
背景
Doris 不是凭空产生数据的。它是个”理解数据”的引擎,数据需要从别处来。
在数据密集型后端架构里,数据来源通常是:
- 业务数据库(MySQL):用户、订单、商品等结构化数据
- 日志/埋点:用户行为、系统日志等半结构化数据
- 外部数据:第三方 API、合作伙伴数据
这些数据到达 Doris 的路径大致是:
1 | 数据源 → 消息队列 → Doris |
消息队列是整个数据管道的”中央动脉”。理解不同的接入方式及其取舍,是搭建数据处理链路的基础。
方式一:Routine Load(最常用)
Routine Load 是 Doris 从 Kafka 持续消费数据的方式。提交一个例行导入任务后,Doris 会持续从 Kafka 拉取数据,实时写入。
1 | CREATE ROUTINE LOAD example_db.routine_load_job ON user_behavior |
优点:持续消费、自动提交 offset、支持 exactly-once(通过 label 机制去重)。
注意:Routine Load 只支持 Kafka。如果你用的是 RocketMQ,需要另外想办法。
方式二:Stream Load(手动推送)
Stream Load 是通过 HTTP 协议把数据推给 Doris。
1 | curl --location-trusted -u user:pass \ |
适用场景:
- 数据是一次性批量导入(不是持续流)
- 从没有 Kafka 的环境中推数据
- 程序内部直接推数据到 Doris(比如 Flink 处理后的结果)
Spring 里可以用 RestTemplate 或 WebClient 调用 Stream Load API。
方式三:Broker Load(批量导入文件)
如果数据已经在 HDFS 或对象存储(S3/OSS)里,用 Broker Load 批量导入。
1 | LOAD LABEL example_db.label_20260721 ( |
适用场景:离线批量导入、数据迁移、定期从 Hive 同步。
方式四:Insert Into(最像 MySQL 的方式)
1 | INSERT INTO user_behavior VALUES |
只适合少量数据。 大批量写入性能很差,因为每条 INSERT 都是一个单独的事务。
RocketMQ 用户怎么办?
Routine Load 不直接支持 RocketMQ。常见方案:
- RocketMQ → Flink → Doris:Flink 消费 RocketMQ,处理后通过 Stream Load 写入 Doris。最灵活,但多一个组件。
- RocketMQ → Kafka Bridge:借助 RocketMQ 的 Kafka 兼容模式或中间桥接服务。简化运维但增加链路。
- 应用层消费 RocketMQ → 批量 Stream Load:在 Spring 应用里消费 RocketMQ 消息,攒够一批再 Stream Load 到 Doris。最简单,适合中小流量。
从 MySQL 到 Doris 的经典链路
最经典的实时同步链路:
1 | MySQL(binlog) → Canal → Kafka → Doris Routine Load |
Canal 伪装成 MySQL 的从库,读取 binlog,把变更事件以 JSON 格式发到 Kafka。Doris Routine Load 从 Kafka 消费这些事件,实时写入。
这套链路的优势是:MySQL 感知不到 Doris 的存在,零侵入。Doris 感知不到 Canal 的存在,只关心 Kafka 里的数据。
但需要注意:
- Canal 的 binlog 格式需要是 ROW(不是 STATEMENT)
- 大事务可能导致 Kafka 消息堆积
- DDL 变更(加列)需要额外处理
选型总结
| 接入方式 | 延迟 | 吞吐 | 运维复杂度 | 适用场景 |
|---|---|---|---|---|
| Routine Load | 秒级 | 高 | 低 | 持续实时导入 |
| Stream Load | 秒级 | 中 | 中 | 程序推送/批量导入 |
| Broker Load | 分钟级 | 最高 | 中 | 离线批量/迁移 |
| Insert Into | 实时 | 极低 | 低 | 少量测试数据 |
一句话:实时场景用 Routine Load(前提是有 Kafka),批量场景用 Broker Load,程序推送用 Stream Load。RocketMQ 用户需要通过 Flink 或应用层桥接。
