Flink 复杂事件处理 CEP
概述
CEP(Complex Event Processing)在事件流中检测复杂模式:从"连续 3 次失败登录"到"下单后 10 分钟未支付"这类业务规则,CEP 用模式匹配实现实时告警与风控。本文讲透 Pattern API、NFA 引擎原理与超时处理。
一、CEP 是什么
1.1 概念
输入:事件流(e1, e2, e3, ...)
模式:事件序列的约束(顺序、数量、时间)
输出:匹配到模式的事件组合| 场景 | 模式示例 |
|---|---|
| 风控 | 1 分钟内 5 次失败登录 |
| 报警 | 连续 3 次超过阈值 |
| 监控 | 下单 30 分钟未支付 |
| 异常检测 | A 后紧跟 B 再 C |
二、Pattern API
2.1 模式定义
java
Pattern<Event, ?> pattern = Pattern
.<Event>begin("start") // 开始
.where(new SimpleCondition<Event>() { // 条件
@Override
public boolean filter(Event event) {
return event.getType() == "login_fail";
}
})
.times(3) // 连续 3 次
.within(Time.minutes(1)); // 1 分钟内2.2 模式操作符
| 操作符 | 说明 |
|---|---|
begin/next/notNext | 序列连接 |
followedBy/notFollowedBy | 非严格连接 |
times/timesOrMore | 次数 |
optional | 可选 |
within | 时间窗口 |
where | 条件 |
or | 条件或 |
2.3 连接方式
next(严格):A 后紧跟 B
followedBy(宽松):A 后任意位置 B
notNext/notFollowedBy:A 后不能有 B| 示例 | 匹配 |
|---|---|
| A.next(B) | A、B 必须相邻 |
| A.followedBy(B) | A 之后出现 B 即可 |
三、模式匹配
3.1 检测执行
java
DataStream<Event> stream = ...;
DataStream<Alert> alerts = CEP
.pattern(stream.keyBy(Event::getUserId), pattern)
.process(new PatternProcessFunction<Event, Alert>() {
@Override
public void processMatch(
Map<String, List<Event>> match, // 匹配事件
Context ctx,
Collector<Alert> out) {
out.collect(buildAlert(match));
}
});3.2 匹配结果
match:模式名 → 匹配事件列表
{"start": [e1, e2, e3]}四、NFA 引擎原理
4.1 什么是 NFA
NFA(Non-deterministic Finite Automaton,非确定有限自动机)是 CEP 的匹配引擎:
模式编译为状态机:
start → login_fail × 3 → 匹配完成
每个状态可接收/拒绝事件
事件流喂入状态机 → 推进/回退匹配4.2 匹配过程
1. 事件到来,从起始状态尝试匹配
2. 满足条件 → 推进状态,记录部分匹配
3. 不满足 → 匹配失败,清理
4. 达到终态 → 输出匹配结果| NFA 特性 | 说明 |
|---|---|
| 非确定性 | 并行追踪多个可能匹配 |
| 状态存储 | 部分匹配保存在状态 |
| 时间窗口 | within 约束淘汰过期 |
4.3 与状态的关系
部分匹配存于 Flink 状态(keyed state)
支持 checkpoint 与精确恢复五、超时处理
5.1 问题
模式设了 within 时间窗口,未在窗口内完成的匹配需要处理。
5.2 超时侧输出
java
OutputTag<Alert> timeoutTag = new OutputTag<Alert>("timeout") {};
DataStream<Alert> result = CEP
.pattern(stream, pattern)
.process(
new PatternProcessFunction<Event, Alert>() {
@Override
public void processMatch(...) { ... } // 正常匹配
},
timeoutTag // 超时输出
);
result.getSideOutput(timeoutTag) // 超时事件流| 处理 | 说明 |
|---|---|
| 正常匹配 | processMatch |
| 超时部分匹配 | 侧输出单独处理 |
| 场景 | 下单未支付提醒 |
5.3 示例:超时订单提醒
模式:下单 → within(30min) → 支付
未支付 → 超时输出 → 发送提醒/自动取消六、业务场景实践
6.1 登录风控
模式:1 分钟内 5 次失败登录(同用户)
输出:锁定账号、告警6.2 指标监控
模式:连续 3 次监控值超阈值
输出:触发告警6.3 业务流程
模式:下单 → 支付成功 → 发货
异常:下单后未支付 → 超时提醒| 场景 | 模式要点 |
|---|---|
| 风控 | keyBy 用户,times + within |
| 告警 | 连续条件 |
| 业务流程 | 序列 + 超时 |
七、性能与限制
| 要点 | 说明 |
|---|---|
| 状态开销 | 部分匹配存状态,注意清理 |
| 并行 | 按 key 并行匹配 |
| 模式复杂度 | 复杂模式状态多 |
| 事件顺序 | 乱序需事件时间与水位线 |
| 优化 | 说明 |
|---|---|
| 精简模式 | 减少状态分支 |
| 状态 TTL | 清理过期匹配 |
| 时间窗口 | within 限制匹配深度 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 匹配不到 | 模式条件/连接方式错误,检查 followedBy |
| 状态增长 | 部分匹配未清理,用 within/TTL |
| 乱序不匹配 | 用事件时间 + 水位线 |
| 超时不触发 | 确认 within 与侧输出 |
| 性能差 | 模式复杂,简化或分区并行 |