F Charge Flink Monitor Guide

Apache Flink Streaming Demo

从充电桩实时监控,理解 Flink 的核心能力

这份 HTML Guide 把理论概念、工程代码和运行现象放在同一条线上:设备遥测进入 Kafka, Flink 用事件时间、Watermark、Keyed State、Timer 和 Window 做实时计算,再写入 PostgreSQL, 最后由 API 和 Dashboard 展示。

1. 学习目标

这个 Demo 不是只演示 API,而是构造了一个真实业务里常见的实时计算闭环。学习时要把 “Flink 概念”放回业务问题里看:乱序事件怎么处理、设备状态怎么记住、离线怎么发现、 分钟级指标怎么聚合、结果怎么被看板消费。

Source

Kafka 输入

模拟器持续写入 pile-telemetry,Flink 使用 KafkaSource 消费。

Time

事件时间

使用设备侧 eventTimeMillis,允许 10 秒乱序。

State

状态与定时器

pileId 维护最后上报时间和离线告警状态。

Sink

结果落库

写入 pile_lateststation_minute_metricsalerts

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
观察顺序 先看 Kafka 是否有数据,再看 Flink 作业是否运行,然后看 PostgreSQL 表是否写入,最后看 Dashboard。 这样排查问题比直接刷新页面更有效。

3. 端到端数据链路

这条链路可以映射到生产系统:充电桩或网关上报遥测,消息队列削峰和解耦,Flink 做实时计算, OLAP/数据库保存最新状态和聚合结果,查询服务给运营大盘和告警系统使用。

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。
生产提醒 PostgreSQL 适合教学、小规模后台和本地闭环。大规模实时分析通常会考虑 ClickHouse、Doris、StarRocks、 Iceberg、Hudi 或专门的告警平台;端到端一致性要结合 Checkpoint、幂等 Sink、事务 Sink 和主键设计一起看。

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 或内存状态。

结果模型

最新状态、明细事件、聚合指标、告警事件最好分开建模;这会直接影响查询性能和后续扩展。

学习闭环 最有效的学习路径是:运行 Dashboard 看现象,打开 Flink UI 看拓扑和 Checkpoint,再去代码里找 Source、State、 Timer、Window、Sink,最后通过一个小规则改造验证自己的理解。