实时数仓架构设计
概述
实时数仓用流处理持续产出低延迟数据:分层思路沿袭离线数仓(ODS/DWD/DWS/ADS),但以 Kafka 为骨架、Flink 为引擎。本文讲清实时分层设计、Kafka 消息标准化,以及与离线数仓的差异。
一、实时数仓是什么
| 对比 | 离线数仓 | 实时数仓 |
|---|---|---|
| 延迟 | 天/小时 | 秒/分钟 |
| 引擎 | Spark/Hive | Flink |
| 存储 | HDFS | Kafka + OLAP |
| 数据 | 批快照 | 流式追加 |
目标:数据秒级可见,支撑大屏、风控、实时报表二、实时分层设计
2.1 分层架构
ODS(实时原始层)
↓ Flink 清洗
DWD(实时明细层)
↓ Flink 聚合
DWS(实时汇总层)
↓
ADS(实时应用层)
DIM(维表):MySQL/Redis/HBase| 层 | 存储 | 内容 |
|---|---|---|
| ODS | Kafka | 原始消息 |
| DWD | Kafka/Iceberg | 清洗明细 |
| DWS | Kafka/OLAP | 分钟级聚合 |
| ADS | OLAP/MySQL | 应用指标 |
2.2 分层职责
| 层 | 职责 |
|---|---|
| ODS | 原样接入、标准化 |
| DWD | 清洗、维表补齐、宽表 |
| DWS | 按维度聚合 |
| ADS | 业务指标输出 |
2.3 分层 vs 直接算
| 方式 | 优点 | 缺点 |
|---|---|---|
| 分层 | 复用、清晰、可回溯 | 链路长延迟高 |
| 直接算 | 快 | 复用差、难维护 |
建议:核心指标分层,简单指标可直算三、Kafka 消息标准化
3.1 为什么标准化
多业务数据进入统一管道:
统一格式 → 统一处理 → 下游复用3.2 标准消息结构
json
{
"event_id": "uuid",
"event_type": "order_created",
"biz_time": 1785918450000,
"data": { "order_id": "1001", "amount": 99.9 },
"trace": { "source": "app", "env": "prod" }
}| 字段 | 说明 |
|---|---|
| event_id | 全局唯一 |
| event_type | 事件类型 |
| biz_time | 业务时间(事件时间) |
| data | 业务数据 |
| trace | 来源/环境 |
3.3 消息格式规范
| 规范 | 说明 |
|---|---|
| 统一 Schema | Avro/JSON Schema |
| 时间戳 | 业务时间必带 |
| 幂等键 | event_id 去重 |
| Topic 规范 | 按主题命名 |
Topic 命名:
ods_xxx(原始)、dwd_xxx(明细)、dws_xxx(汇总)四、链路设计
4.1 典型链路
业务库 MySQL
→ Flink CDC → Kafka(ods)
→ Flink 清洗 → Kafka(dwd) / Iceberg
→ Flink 聚合 → Kafka(dws) / Doris
→ 大屏/报表/风控| 环节 | 组件 |
|---|---|
| 采集 | Flink CDC / SDK 埋点 |
| 管道 | Kafka |
| 加工 | Flink SQL |
| 存储 | Kafka、Iceberg、Doris/ClickHouse |
| 服务 | 查询接口/大屏 |
4.2 双链路(批流一体)
实时:Kafka → Flink → Doris(秒级)
离线:HDFS → Spark → 数仓(天级)
↓ 结果对账、互为校验五、与离线数仓的差异
| 维度 | 离线 | 实时 |
|---|---|---|
| 计算模型 | 批 | 流(微批/持续) |
| 数据存储 | Hive/HDFS | Kafka/OLAP |
| 时间语义 | 分区日期 | 事件时间/水位线 |
| 准确性 | 重算校正 | 依赖流一致性 |
| 回溯 | 简单重跑 | 需重放/补偿 |
5.1 准确性差异
| 实时痛点 | 处理 |
|---|---|
| 迟到数据 | 水位线 + 延迟允许 |
| 重复数据 | 幂等/去重 |
| 指标不一致 | 与离线对账 |
5.2 融合趋势
湖仓一体(Iceberg/Paimon):
实时写湖 → 批读湖 → 一套数据
流批统一,减少双链路六、架构选型
6.1 组件选型
| 环节 | 选项 |
|---|---|
| 采集 | Flink CDC、Canal、SDK |
| 管道 | Kafka(标准) |
| 加工 | Flink SQL(主流) |
| 明细存储 | Kafka、Iceberg、Paimon |
| 汇总存储 | Doris、ClickHouse |
| 维表 | MySQL、Redis、HBase |
6.2 简单 vs 复杂
指标少、延迟要求高 → 简化分层(ODS → ADS)
指标多、复用需求 → 完整分层七、设计要点
| 要点 | 说明 |
|---|---|
| 事件时间贯穿 | 统一业务时间 |
| 消息标准化 | 管道可复用 |
| 幂等设计 | 防重复 |
| 水位线统一 | 全局时间语义 |
| 对账机制 | 与离线互验 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 实时离线不一致 | 口径差异,统一指标定义 |
| 消息格式混乱 | 标准化 + Schema 管理 |
| 链路太长延迟高 | 简化分层 |
| 迟到数据影响指标 | 水位线 + 修正机制 |
| 重复计算 | 幂等 + 去重 |