数据管道 Data Pipeline
概述
一条完整的数据管道把数据从采集端送到分析端:Kafka → Flink → Iceberg → Doris/ClickHouse。本文以电商实时指标为例,讲清管道分层、各环节配置、容错与监控。
一、管道全景
业务日志/业务库
→ Kafka(ODS 消息)
→ Flink(实时加工)
→ Iceberg(明细落湖)
→ Doris/ClickHouse(指标查询)
→ 大屏/报表/API| 环节 | 组件 | 职责 |
|---|---|---|
| 采集 | Flink CDC / 埋点 | 进管道 |
| 管道 | Kafka | 缓冲与分发 |
| 加工 | Flink SQL | 清洗聚合 |
| 明细 | Iceberg | 可回溯 |
| 汇总 | Doris/CH | 秒级查询 |
二、Kafka 环节
2.1 Topic 设计
| Topic | 内容 |
|---|---|
| ods_orders | 原始订单 |
| dwd_orders | 加工明细 |
| dws_order_min | 分钟聚合 |
| ads_order_realtime | 应用指标 |
分区:
按业务量规划,支持并行消费
关键 topic 按 key 分区保序2.2 参数
| 参数 | 说明 |
|---|---|
| 副本数 | 生产 3 副本 |
| 保留时间 | 按回溯需求 |
| 压缩 | 可选 |
三、Flink 加工环节
3.1 作业拆分
| 作业 | 逻辑 |
|---|---|
| CDC → Kafka | 采集作业 |
| Kafka → DWD | 清洗作业 |
| DWD → DWS | 聚合作业 |
| DWS → Doris | 输出作业 |
按职责拆分作业:
独立重启、独立调优
避免单作业过大3.2 Flink SQL 示例
sql
-- DWD:清洗
CREATE VIEW dwd_orders AS
SELECT order_id, user_id, amount, biz_time
FROM kafka_ods_orders
WHERE amount > 0;
-- DWS:分钟聚合
INSERT INTO kafka_dws
SELECT TUMBLE_START(biz_time, INTERVAL '1' MINUTE) AS win_start,
COUNT(*) AS cnt, SUM(amount) AS gmv
FROM dwd_orders
GROUP BY TUMBLE(biz_time, INTERVAL '1' MINUTE);3.3 配置要点
| 配置 | 说明 |
|---|---|
| Checkpoint | 60-120s,精确一次 |
| 并行度 | 与 Kafka 分区匹配 |
| 水位线 | 事件时间贯穿 |
| 状态 TTL | 防膨胀 |
四、Iceberg 落湖环节
4.1 落湖目的
| 目的 | 说明 |
|---|---|
| 明细留存 | 可回溯/重算 |
| 湖仓一体 | 离线批读 |
| 对账 | 与实时结果比对 |
4.2 写入
sql
CREATE TABLE lake.dwd_orders (...) WITH (
'connector' = 'iceberg',
'catalog-name' = 'lake',
'catalog-type' = 'hive',
'uri' = 'thrift://metastore:9083',
'write.format.default' = 'parquet'
);
INSERT INTO lake.dwd_orders SELECT * FROM kafka_dwd_orders;4.3 落湖策略
| 策略 | 说明 |
|---|---|
| 全量落 | 全部明细 |
| 采样落 | 部分明细 |
| 按层落 | 核心层 |
落湖与实时并行:
实时算指标,湖存明细
两者数据同源五、Doris / ClickHouse 环节
5.1 选型
| 引擎 | 特点 | 适用 |
|---|---|---|
| Doris | 高并发、多表 Join | 报表/查询 |
| ClickHouse | 极快聚合 | 大屏/聚合 |
5.2 写入方式
sql
-- Doris Sink
CREATE TABLE doris_sink (...) WITH (
'connector' = 'doris',
'fenodes' = 'doris:8030',
'table.identifier' = 'dws.order_min',
'sink.label-prefix' = 'flink'
);| 要点 | 说明 |
|---|---|
| 幂等 | label 保证不重 |
| 批量 | 批量写入 |
| 主键 | 聚合表模型 |
5.3 查询
sql
SELECT city, SUM(gmv) FROM dws.order_min
WHERE win_start >= '2026-08-04 10:00:00'
GROUP BY city;六、管道容错
| 环节 | 容错 |
|---|---|
| Kafka | 副本 + 保留重放 |
| Flink | Checkpoint 恢复 |
| Iceberg | 快照/事务 |
| Doris | 幂等写入 |
端到端一致性:
Source offset + 状态 + 幂等 Sink
→ Exactly-Once6.1 故障处理
| 故障 | 处理 |
|---|---|
| Flink 作业失败 | 重启 + checkpoint 恢复 |
| Kafka 抖动 | 重试/背压 |
| 下游慢 | 削峰/扩容 |
| 数据异常 | 侧输出 + 修复 |
七、管道监控
7.1 指标
| 指标 | 说明 |
|---|---|
| 端到端延迟 | 事件到查询 |
| 吞吐 | 各环节速率 |
| 积压 | Kafka 未消费 |
| 失败率 | 作业/写入 |
7.2 监控看板
Prometheus + Grafana:
Flink JobManager/TaskManager 指标
Kafka Lag
写入速率与延迟
作业状态(失败/重启)八、管道优化
| 优化 | 说明 |
|---|---|
| 并行匹配 | 各环节并行度匹配 |
| 减少拷贝 | 避免重复计算 |
| 合并小写 | 批量写入 |
| 分区对齐 | Kafka 分区与并行 |
| 网络 | 同机房部署 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 端到端延迟高 | 定位瓶颈环节 |
| Kafka 积压 | 扩容消费者 |
| 明细丢失 | Checkpoint + 幂等 |
| 查询慢 | Doris 模型/索引优化 |
| 数据不一致 | 对账 + 修正 |