实时 ETL 设计
概述
实时 ETL 是实时数仓的核心加工环节:把原始数据清洗、转换、补齐、关联成可分析数据。本文讲清 CDC 入湖、清洗转换、维表关联、双流 Join,以及实时数据质量保障。
一、实时 ETL 全景
原始消息(Kafka ODS)
→ 清洗(过滤/纠正/标准化)
→ 转换(类型/脱敏/补字段)
→ 补齐(维表关联)
→ 关联(双流 Join)
→ 明细(Kafka DWD / 湖)| 环节 | 职责 |
|---|---|
| 清洗 | 去脏、纠正、过滤 |
| 转换 | 格式、类型、脱敏 |
| 补齐 | 维表字段 |
| Join | 多流关联 |
二、CDC 入湖
2.1 链路
MySQL binlog → Flink CDC → Kafka(ODS)→ 湖/下游2.2 入湖方式
| 方式 | 说明 |
|---|---|
| Kafka → Iceberg/Paimon | 明细落湖 |
| Kafka → 下游 Kafka | 继续加工 |
| 直接算 | 不入湖 |
2.3 注意
| 要点 | 说明 |
|---|---|
| 主键 | 变更按主键 upsert |
| Schema 变更 | 评估字段影响 |
| 位点 | Checkpoint 保证不丢 |
三、清洗与转换
3.1 清洗
| 清洗 | 示例 |
|---|---|
| 过滤 | 丢弃空数据/测试数据 |
| 纠正 | 修复格式错误 |
| 去重 | event_id 去重 |
| 标准化 | 统一枚举/单位 |
sql
-- 过滤脏数据
CREATE VIEW clean AS
SELECT * FROM ods_orders
WHERE amount > 0 AND user_id IS NOT NULL;3.2 转换
| 转换 | 示例 |
|---|---|
| 类型 | 字符串 → 数值 |
| 时间 | 时间戳格式化 |
| 脱敏 | 手机号打码 |
| 补字段 | 添加分区键、标签 |
3.3 时间标准化
统一事件时间:
消息 biz_time → 标准时间戳
用于水位线与窗口四、维表关联
4.1 方案
| 方案 | 说明 | 适用 |
|---|---|---|
| Lookup Join | 实时查维表 | 小维表 |
| 广播维表 | 内存加载 | 高频查询 |
| 维表流 | CDC 维表转流 | 变化频繁 |
4.2 Lookup Join
sql
SELECT o.order_id, u.city, u.level
FROM orders o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
ON o.user_id = u.user_id;4.3 维表缓存
| 参数 | 说明 |
|---|---|
lookup.cache.max-rows | 缓存行数 |
lookup.cache.ttl | 缓存时间 |
缓存提升性能,容忍短时陈旧4.4 维表更新
维表数据变化 → 重建/缓存过期
大维表 → HBase/Redis 存储五、双流 Join
5.1 场景
订单流 × 支付流 → 关联订单支付
点击流 × 曝光流 → 关联点击转化5.2 双流 Join 方式
sql
SELECT ... FROM orders o
INNER JOIN payments p
ON o.order_id = p.order_id
AND o.ts BETWEEN p.ts AND p.ts + INTERVAL '5' MINUTE;5.3 双流 Join 关键点
| 要点 | 说明 |
|---|---|
| 时间区间 | 关联窗口约束 |
| 状态 | 保留未匹配数据 |
| 状态清理 | TTL 清理过期 |
| 数据补齐 | 迟到匹配 |
5.4 三种 Join
| 类型 | 语义 |
|---|---|
| INNER | 双流都到才输出 |
| LEFT | 左流先输出,右流补齐 |
| FULL | 双流补全 |
Interval Join:
关联时间窗口内的数据
状态大小可控六、数据质量保障
6.1 实时质量检查
| 检查 | 说明 |
|---|---|
| 空值率 | 关键字段缺失 |
| 异常值 | 数值范围 |
| 迟到率 | 超过水位线比例 |
| 流量突降 | 采集中断信号 |
6.2 质量看板
指标:
处理条数、成功/失败率
迟到率、去重率
维表命中率6.3 质量兜底
| 兜底 | 说明 |
|---|---|
| 侧输出 | 脏数据单独处理 |
| 告警 | 异常触发 |
| 重放 | 从 Kafka 重放修正 |
七、实时 ETL 最佳实践
| 实践 | 说明 |
|---|---|
| 幂等 | 结果可重放 |
| 事件时间 | 贯穿处理 |
| 状态清理 | 防无限增长 |
| 并行度 | 与吞吐匹配 |
| 监控 | 延迟/质量指标 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 维表命中低 | 维表未同步,缓存过期 |
| 双流 Join 状态大 | 时间区间过大,加 TTL |
| 数据重复 | 去重键缺失 |
| 脏数据流入 | 清洗规则 + 侧输出 |
| 质量异常不感知 | 建立实时质量监控 |