基于 Flink + Kafka + ClickHouse 的电商用户行为实时分析系统,涵盖数据采集、流计算、多层数仓、可视化全链路。
数据源 ──→ Kafka ──→ Flink ──→ ClickHouse(ODS/DWD/DWS) ──→ Spring Boot API ──→ ECharts 仪表盘
↑ │
├── BehaviorProducer(模拟) ├── 字段校验 → 脏数据异常表
└── LogFileProducer(日志) ├── 告警检测 → alert_log
└── UV 去重(HashSet + HyperLogLog)
| 层级 | 技术 |
|---|---|
| 数据采集 | Java Kafka Producer / CSV 日志回放 |
| 消息队列 | Kafka 3.7 (KRaft) |
| 流计算 | Apache Flink 1.17 (Checkpoint + Watermark + Window) |
| 存储 | ClickHouse 24.3 (ODS → DWD → DWS 三层) |
| 后端 | Spring Boot 2.7 |
| 前端 | ECharts 5.5 深色仪表盘 |
- UV 去重:混合策略 — HashSet 精确(< 1万) + HyperLogLog 近似(> 1万),含误差推导
- 状态后端:生产级 RocksDB 增量 Checkpoint(代码标注,Windows 开发用 HashMap)
- 数仓分层:ODS(原始) → DWD(明细) → DWS(汇总),ReplacingMergeTree 实现幂等
- 数据质量:字段级校验(非空/格式/值域) + 脏数据异常表 + PV/转化率告警
- 容错:Checkpoint 10s + Sink 3次重试 + 后端 503 + Wm+allowedLateness+侧输出
- JDK 11+
- Maven 3.9+
- Docker Desktop
docker compose -f docker/docker-compose.yml up -dcurl -X POST "http://localhost:8123/" --data-binary "
CREATE TABLE IF NOT EXISTS metrics_realtime (
window_end DateTime, pv UInt64, uv UInt64,
view_count UInt64, cart_count UInt64, buy_count UInt64, pay_count UInt64,
view_to_cart Float64, cart_to_buy Float64, buy_to_pay Float64,
ck_version UInt64
) ENGINE = ReplacingMergeTree(ck_version) ORDER BY window_end;
"IDEA 中依次运行:
BehaviorProducer.main()— 数据生产UvPvFunnelAnalysis.main()— Flink 计算Application.main()— 后端 + 前端
浏览器打开 http://localhost:8080/
├── docker/docker-compose.yml # Kafka + ClickHouse
├── data/behavior_log.csv # 模拟日志数据
├── user-behavior-analysis/
│ └── src/main/
│ ├── java/org/example/
│ │ ├── kafka/ # 数据生产者
│ │ ├── flink/ # Flink 流计算
│ │ ├── model/ # 数据模型
│ │ └── controller/ # REST API
│ └── resources/
│ ├── static/index.html # 仪表盘前端
│ └── application.yml
| 接口 | 说明 |
|---|---|
GET /api/metrics?limit=N |
最近 N 条完整指标 |
GET /api/pv/latest |
最新 PV |
GET /api/uv/latest |
最新 UV |
GET /api/funnel/latest |
最新转化漏斗 |