一、问题场景
一条产线,每秒上千条传感器数据
某电机产线,每台设备有温度、振动、电流 3 路传感器,采样频率 10Hz。100 台设备并发 →每秒约 3000 条遥测。传统做法是"存下来再分析",等发现轴承过热,产线可能已经烧了。
我们要的闭环是:边缘采样 → Kafka 汇聚 → Flink 实时算异常 → 云端告警/看板 → 反向下发降速指令。本篇聚焦中间那段"实时异常检测",也是周五连载云边协同的"大脑"部分。
二、方案设计
整体数据流:
[边缘网关] --MQTT--> [Kafka topic: sensor.raw] | [Flink Job] keyBy(deviceId) → 滑动窗口(z-score) → 异常判定 → ├─ 正常 → 写入时序库(Put) └─ 异常 → 告警(WebSocket/邮件) + 标记为什么用z-score + 滑动窗口而不是简单阈值?因为同一台电机在不同工况下"正常温度"不一样,绝对阈值会误报。用最近窗口的均值/标准差做动态基线,更鲁棒。
三、分步实现(PyFlink,可读性优先)
1. 定义数据结构与源
from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.connectors.kafka import KafkaSource, KafkaOffsets from pyflink.common.serialization import SimpleStringSchema from pyflink.common.watermark_strategy import WatermarkStrategy import json env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(4) source = KafkaSource.builder() \ .set_bootstrap_servers("kafka:9092") \ .set_topics("sensor.raw") \ .set_group_id("flink-anomaly") \ .set_starting_offsets(KafkaOffsets.latest()) \ .set_value_only_deserializer(SimpleStringSchema()) \ .build() ds = env.from_source(source, WatermarkStrategy.no_watermarks(), "kafka")2. 解析 + keyBy 设备
def parse(record): e = json.loads(record) return (e["deviceId"], e["metric"], float(e["value"]), int(e["ts"])) parsed = ds.map(parse, output_type=...) keyed = parsed.key_by(lambda x: (x[0], x[1])) # 按 设备+指标 分组3. 滑动窗口 z-score 异常检测(核心算子)
from pyflink.datastream.window import SlidingEventTimeWindows from pyflink.common.time import Time windowed = keyed \ .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.seconds(30))) \ .process(AnomalyDetector()) class AnomalyDetector(KeyedProcessWindowFunction): def process(self, key, ctx, events): vals = sorted([e[2] for e in events]) n = len(vals) mean = sum(vals) / n var = sum((v - mean) ** 2 for v in vals) / n std = var ** 0.5 + 1e-6 # 用窗口末尾的点做 z-score 判定 latest = vals[-1] z = (latest - mean) / std if abs(z) > 3.0: # 3σ 准则 yield { "deviceId": key[0], "metric": key[1], "value": latest, "z": round(z, 2), "mean": round(mean, 2), "ts": ctx.current_watermark }4. 异常分流到告警 Sink
anomalies = windowed.map(lambda a: json.dumps(a)) anomalies.add_sink(KafkaSink.builder() .set_bootstrap_servers("kafka:9092") .set_record_serializer(..., topic="sensor.alert") .build())下游一个 Spring Boot / Node 服务订阅sensor.alert,推 WebSocket 到运维看板,并按设备 ID 触发降级指令回写边缘网关。
四、踩坑记录
乱序事件必须有 Watermark:工业网关网络抖动,事件迟到是常态。不设 watermark + 允许的延迟,窗口会提前触发导致漏检。
状态膨胀:
keyBy(deviceId, metric)后窗口状态随时间增长,务必配State TTL,否则一周后 JobManager 内存爆炸。z-score 对突发不敏感:纯统计方法抓不出"缓变劣化"。生产里常叠加斜率检测 / EWMA,本篇留给进阶版。
不要在 process 里查数据库:每条事件去查设备元数据会拖垮吞吐,预先广播(
BroadcastState)下发设备配置。
五、性能数据(单机基准)
| 指标 | 数值 |
|---|---|
| 吞吐 | 单 TaskManager(4 核)约12 万 events/s |
| 端到端延迟(采样→告警) | p99 <800ms |
| 100 台设备 3 路传感器 | 稳态 CPU ~55% |
Flink 把"事后看报表"变成了"事中拦风险"。这套管道正是我们整个 Edge AI 全栈的数据主动脉——边缘负责采和跑轻模型,云端 Flink 负责 aggregation 和全局异常判定。