1. 学习目标
这个 Demo 不是只演示 API,而是构造了一个真实业务里常见的实时计算闭环。学习时要把 “Flink 概念”放回业务问题里看:乱序事件怎么处理、设备状态怎么记住、离线怎么发现、 分钟级指标怎么聚合、结果怎么被看板消费。
Kafka 输入
模拟器持续写入 pile-telemetry,Flink 使用 KafkaSource 消费。
事件时间
使用设备侧 eventTimeMillis,允许 10 秒乱序。
状态与定时器
按 pileId 维护最后上报时间和离线告警状态。
结果落库
写入 pile_latest、station_minute_metrics、alerts。
2. 快速启动
在仓库根目录执行:
docker compose up --build -d
| 入口 | 地址 | 你应该观察什么 |
|---|---|---|
| 实时大盘 | http://localhost:3000 |
KPI、站点趋势、告警、充电桩状态墙是否刷新。 |
| Flink UI | http://localhost:8082 |
作业 Charge Pile Realtime Monitoring、算子拓扑、Checkpoint。 |
| Kafka UI | http://localhost:8081 |
pile-telemetry topic 中是否持续产生 JSON 消息。 |
| API Docs | http://localhost:8000/docs |
查询 /api/kpis、/api/alerts、/api/power-trend。 |
3. 端到端数据链路
这条链路可以映射到生产系统:充电桩或网关上报遥测,消息队列削峰和解耦,Flink 做实时计算, OLAP/数据库保存最新状态和聚合结果,查询服务给运营大盘和告警系统使用。
4. Flink 理论到代码
入口文件是 flink-job/src/main/java/com/example/charge/ChargingPileMonitoringJob.java。
Source:从 Kafka 读 JSON
生产系统里可替换成 MQTT 网关、业务事件总线或 CDC 后的 Kafka topic。
KafkaSource<PileTelemetry> source =
KafkaSource.<PileTelemetry>builder()
.setTopics(topic)
.setGroupId("charge-monitoring-flink")
.setValueOnlyDeserializer(
new PileTelemetryDeserializationSchema())
.build();
Event Time:用设备事件时间
模拟器会制造 5% 的乱序事件,Guide 中的 Watermark 正是为这个问题服务。
WatermarkStrategy
.<PileTelemetry>forBoundedOutOfOrderness(
Duration.ofSeconds(10))
.withTimestampAssigner(
(event, ts) -> event.eventTimeMillis)
.withIdleness(Duration.ofSeconds(30));
5. Keyed State 与 Timer:离线告警为什么能做准
离线不是单条事件能判断出来的,它需要记住“每个桩最后一次上报是什么时候”。
因此 Flink 先按 pileId 分区,再在每个 Key 下保存状态。
lastSeenTs
保存当前 pileId 最近一次事件时间,用于判断 Timer 是否过期。
lastStationId
Timer 触发时仍能知道离线桩属于哪个站点。
offlineAlertOpen
避免同一个离线周期重复输出多条离线告警。
6. Window 聚合:从单桩遥测到站点指标
单条遥测只说明某个桩的瞬时状态,运营大盘更关心“每个站点最近一分钟的负载和异常”。
这就是 keyBy(stationId) 加滚动事件时间窗口的价值。
events
.keyBy(event -> event.stationId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new StationMetricAggregate(), new StationMetricWindowFunction())
7. Sink 与查询层:结果如何被业务使用
Demo 使用 PostgreSQL 方便教学和本地演示。Flink 输出三类结果,API 再把它们组合成 Dashboard 能直接消费的查询接口。
| 表 | 写入方式 | 业务含义 | Dashboard 用途 |
|---|---|---|---|
pile_latest |
按 pile_id upsert |
每个充电桩最新状态。 | KPI、状态墙、站点当前分布。 |
station_minute_metrics |
按 window_start + station_id upsert |
站点分钟级窗口指标。 | 功率趋势、站点榜单、负载分析。 |
alerts |
按 alert_id 去重插入 |
故障、过温、功率尖峰、离线、恢复。 | 最近告警列表、10 分钟告警 KPI。 |
8. 工程代码地图
学习时建议从运行现象反推代码,不要从 Maven 项目结构开始硬读。
| 学习问题 | 看哪个文件 | 关注点 |
|---|---|---|
| 数据是怎么产生的? | simulator/simulator.py |
站点模板、状态切换、乱序事件、故障和离线模拟。 |
| Flink 作业主链路在哪里? | ChargingPileMonitoringJob.java |
KafkaSource、Watermark、告警流、窗口聚合流、JDBC Sink。 |
| 离线告警怎么实现? | AlertDetector.java |
ValueState、registerEventTimeTimer、onTimer、恢复告警。 |
| 站点聚合怎么计算? | StationMetricAggregate.java |
累计功率、温度、桩数和状态计数。 |
| 窗口补充字段在哪里? | StationMetricWindowFunction.java |
写入窗口开始和结束时间、stationId。 |
| 结果落库 SQL 在哪里? | PostgresSinks.java |
upsert、insert、主键冲突处理。 |
| API 怎么查大盘数据? | api/main.py |
KPI、站点、充电桩、告警、趋势接口。 |
9. 动手实验:从看懂到改动
把离线阈值改短
修改 docker-compose.yml 中 Flink submit 的 OFFLINE_TIMEOUT_MS,例如改成 15000,再重启作业。
docker compose up --build -d flink-submit
新增“站点故障风暴”规则
目标:同一站点 5 分钟内故障事件超过 5 个时输出 STATION_FAULT_STORM。
events
.filter(PileTelemetry::hasFault)
.keyBy(event -> event.stationId)
.window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1)))
把站点聚合改成 Flink SQL
用 SQL 重写分钟窗口,可以练习 Table API 和 DataStream API 的边界。
SELECT
stationId,
TUMBLE_START(rowtime, INTERVAL '1' MINUTE) AS window_start,
COUNT(*) AS total_events
FROM telemetry
GROUP BY stationId, TUMBLE(rowtime, INTERVAL '1' MINUTE)
加入站点维表
增加 station_dim,把运营商、省份、额定容量、价格策略等维度关联到实时指标里。
station_dim(station_id, operator_id, province, capacity_kw, price_policy_id)
10. 生产化思考
可靠性
关注 Checkpoint 间隔、状态后端、Kafka offset 提交、Sink 幂等和重启策略。Demo 已启用
EXACTLY_ONCE checkpoint,但端到端语义仍取决于 Sink。
时间语义
Watermark 延迟决定窗口完整性与结果时效之间的取舍。IoT 场景常见设备补发、网关缓存和网络抖动。
状态规模
设备数变大后,Keyed State 会成为核心资源。需要设置 TTL、规划 key 分布,并监控 RocksDB 或内存状态。
结果模型
最新状态、明细事件、聚合指标、告警事件最好分开建模;这会直接影响查询性能和后续扩展。