摘要:会写 Pattern 不等于会开发 CEP 作业——从规则到落地的完整链路是"模式定义 → 检测 → 选择"三阶段。这篇文章按生产视角拆解:定义阶段的条件编写(SimpleCondition 与 IterativeCondition 的差别,后者能引用同模式已匹配事件实现"连续递增"这类跨事件条件)、检测阶段的准备与绑定(keyBy、事件时间、CEP.pattern 的 NFA 运行语义)、选择阶段的三档 API(select/flatSelect/process 怎么选、超时出口怎么写、Map<String, List> 怎么取数)。四个完整代码案例覆盖从简单告警到迭代条件、多路输出与超时处理,读完能独立搭出一个生产级 CEP 作业。
关键词:Flink CEP、模式定义、SimpleCondition、IterativeCondition、检测、PatternStream、NFA、选择、PatternSelectFunction、PatternFlatSelectFunction、PatternProcessFunction、TimedOutPartialMatchHandler、代码实现
一、CEP 开发的完整链路:不是"写个 Pattern 就行"
CEP 篇讲了 NFA 引擎,Pattern 篇讲了语法细节。但真实开发里,一个 CEP 作业的代码组织是另一回事——它分三个阶段:
- 模式定义:把业务规则翻译成 Pattern(怎么写条件、怎么命名);
- 检测:把 Pattern 挂到流上让引擎跑起来(怎么准备流、怎么绑定);
- 选择:从匹配结果里提取业务数据(怎么选 API、怎么处理超时)。
三阶段各有各的 API 和坑。这篇按生产视角走一遍完整链路。
二、三阶段全景
| 阶段 | 输入 | 输出 | 关键 API |
|---|---|---|---|
| 定义 | 业务规则 | Pattern 对象 | begin/where/next/times/within |
| 检测 | Pattern + 事件流 | PatternStream | keyBy + CEP.pattern() |
| 选择 | PatternStream | DataStream + 侧输出 | select/flatSelect/process |
定义阶段是纯声明(不执行匹配,只描述规则),检测阶段是引擎执行(用户只看到 PatternStream),选择阶段是结果落地(匹配与超时两条出口都在这)。
三、模式定义:条件编写的两种姿势
3.1 SimpleCondition:单事件条件
最简单也最常用,只根据当前事件判断:
Pattern<Transaction,Transaction>p=Pattern.<Transaction>begin("large").where(newSimpleCondition<Transaction>(){@Overridepublicbooleanfilter(Transactiont){returnt.amount>100_000;// 只看当前事件,够用}});3.2 IterativeCondition:迭代条件(能引用已匹配事件)🔥
需求一复杂就发现 SimpleCondition 不够用了——比如反洗钱的经典模式"连续多笔交易金额逐笔递增":判断当前事件"金额比上一笔大 2 倍"必须知道上一笔是多少。这时用IterativeCondition,它通过ctx.getEventsForPattern("模式名")拿到同模式内之前已匹配的事件:
Pattern<Transaction,Transaction>escalating=Pattern.<Transaction>begin("tx").where(newIterativeCondition<Transaction>(){@Overridepublicbooleanfilter(Transactioncurrent,Context<Transaction>ctx){// 第一笔:无条件进入if(!ctx.getEventsForPattern("tx").iterator().hasNext()){returncurrent.amount>10_000;// 起点:大额}// 后续事件:金额必须比上一笔大 2 倍(跨事件条件!)Transactionprev=ctx.getEventsForPattern("tx").iterator().next();// 同模式已匹配的上一笔returncurrent.amount>prev.amount*2;}}).times(3)// 3 笔逐笔翻倍.within(Time.minutes(5));这是 CEP 表达力的分水岭:SimpleCondition 看单事件,IterativeCondition 看事件序列的上下文。“连续递增”“比上一笔大”"首笔触发后行为变化"这类规则,只有 IterativeCondition 能优雅表达——用 SimpleCondition 只能靠外部状态 hack。
3.3 命名规范:模式名 = 选择阶段的取数 key
begin("start")里的名字不是装饰——它是选择阶段match.get("start")的 key。规范:语义化、唯一、与选择函数里的引用完全一致。拼错一个字,编译期不报错,运行期 NPE。
四、检测:把 Pattern 挂到流上
检测阶段只有三步,但前两步决定一切:
// ① 准备:keyBy + 事件时间(缺一不可)KeyedStream<Transaction,String>keyed=txStream.keyBy(Transaction::getAccountId)// 每账户独立检测.assignTimestampsAndWatermarks(WatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(5)).withTimestampAssigner((t,ts)->t.getEventTs()));// ② 绑定:流 + 模式 + 跳过策略 → PatternStreamPatternStream<Transaction>ps=CEP.pattern(keyed,escalating,AfterMatchSkipStrategy.skipPastLastEvent());// ③ 产物:PatternStream 已包含匹配 + 超时两类记录,等待选择阶段消费检测的底层就是 NFA 运行(CEP 篇讲过):每个 key 一个独立 NFA 实例,事件按 key 路由、驱动状态转移;within 定时器是状态,随 checkpoint 恢复。这些对用户透明——你只看到 PatternStream。但两个准备不做,检测就是错的:不 keyBy,跨账户的事件会串进同一个匹配;不配事件时间,within 超时永不触发。
五、选择:三档 API 与超时出口
5.1 三档 API 怎么选
| API | 一条匹配输出 | 超时写法 | 上下文 | 适用 |
|---|---|---|---|---|
| select | 1 条 | select(tag, timeoutFn, selectFn) | 无 | 简单告警 |
| flatSelect | 多条 | flatSelect(tag, flatTimeoutFn, flatSelectFn) | 无 | 告警+指标 |
| process | 多条 | 实现 TimedOutPartialMatchHandler | 完整 | 生产主力 |
三档都接收Map<String, List<IN>>——按模式名取该模式匹配到的事件列表:match.get("order")取下单事件,match.get("pay")取支付事件;times 循环时列表里有多个元素;optional 模式可能缺失要判空;超时部分匹配的 Map 只有已匹配的模式(比如只有 “order” 没有 “pay”)。
5.2 完整案例一:三阶段串起来的简单告警(select)
// 需求:5 分钟内连续 3 次登录失败 → 告警(最简单链路)// ── 定义 ──Pattern<LoginEvent,LoginEvent>p=Pattern.<LoginEvent>begin("start").where(e->e.result==FAIL).next("mid").where(e->e.result==FAIL).times(2).consecutive().within(Time.minutes(5));// ── 检测 ──PatternStream<LoginEvent>ps=CEP.pattern(logins.keyBy(LoginEvent::getUserId),p);// ── 选择 ──DataStream<Alert>alerts=ps.select((Map<String,List<LoginEvent>>match)->{LoginEventfirst=match.get("start").get(0);// 按模式名取事件returnnewAlert(first.userId,"brute-force");});5.3 完整案例二:一次命中多路输出(flatSelect)
命中一条撞库模式,同时落告警和监控指标两条记录:
DataStream<Object>result=ps.flatSelect((Map<String,List<LoginEvent>>match,Collector<Object>out)->{LoginEventfirst=match.get("start").get(0);out.collect(newAlert(first.userId,"brute-force"));// 告警out.collect(newMetric("brute-force",1));// 指标// 一次匹配 → 两条输出,下游各自消费});5.4 完整案例三:process + 超时处理(生产主力)
订单超时场景,匹配与超时在一个类里收口,超时走侧输出:
// 模式:下单后 10 分钟内未支付 → 超时Pattern<OrderEvent,OrderEvent>p=Pattern.<OrderEvent>begin("order").where(e->e.type==CREATE).followedBy("pay").where(e->e.type==PAY).within(Time.minutes(10));PatternStream<OrderEvent>ps=CEP.pattern(orders.keyBy(OrderEvent::getOrderId),p);OutputTag<OrderEvent>timeoutTag=newOutputTag<OrderEvent>("timeout"){};DataStream<String>result=ps.process(newPatternProcessFunction<OrderEvent,String>(){@OverridepublicvoidprocessMatch(Map<String,List<OrderEvent>>match,Contextctx,Collector<String>out){// 完整匹配:下单 → 支付out.collect("PAID:"+match.get("pay").get(0).getOrderId());}@OverridepublicvoidhandleTimeout(Map<String,List<OrderEvent>>partial,longts,Contextctx)throwsException{// 超时部分匹配:只有 order(pay 模式缺失)// 注意这里也能侧输出——超时与匹配两条链路一个类收口ctx.output(timeoutTag,partial.get("order").get(0));}});DataStream<OrderEvent>timeoutOrders=result.getSideOutput(timeoutTag);// timeoutOrders → 关单/提醒;result → 正常支付订单process 相比 select 的优势在代码组织上最明显:processMatch 和 handleTimeout 是同一个类的两个方法,超时逻辑就在匹配逻辑旁边,不用像 select 那样把超时函数拆到另一个匿名类里。生产环境我默认选 process。
5.5 迭代条件 + 选择的组合:逐笔递增交易告警
把第三节的 IterativeCondition 模式接上选择,就是一个完整的反洗钱检测:
// 定义(迭代条件,见 3.2 的 escalating 模式)+ 检测 + 选择DataStream<RiskAlert>alerts=CEP.pattern(keyedTx,escalating).flatSelect((Map<String,List<Transaction>>match,Collector<RiskAlert>out)->{// times(3) 循环:match.get("tx") 里有 3 笔事件List<Transaction>txs=match.get("tx");out.collect(newRiskAlert(txs.get(0).accountId,"escalating",txs));});六、实战避坑清单
- 迭代条件别滥用:IterativeCondition 每次评估都要遍历已匹配事件,大循环模式 + 高频事件会放大开销——只在确实需要跨事件条件时用;
- match.get() 前先想 optional:optional 模式可能缺失,直接 get().get(0) 会 NPE;
- 超时 Map 缺模式:handleTimeout 里的 Map 只有已匹配的模式,别拿"还没到"的模式取数;
- 模式名一致:定义与选择两处引用必须完全一致(建议抽 static final 常量);
- 检测前准备两行不能省:keyBy + assignTimestampsAndWatermarks,漏一个语义就错;
- 选择函数别做重活:每条匹配调用一次,里面查库/调用外部系统会拖慢整个 PatternStream;
- process 优先:新代码默认 PatternProcessFunction,模板代码和 select 差不多,但超时与上下文能力完整。
七、总结:我的判断
CEP 开发的完整链路可以用一句话概括:定义是规则、检测是引擎、选择是出口。三个阶段的关注点完全不同——定义阶段想清楚"条件怎么写"(Simple vs Iterative),检测阶段做好"两个准备"(keyBy + 事件时间),选择阶段定好"出口怎么落"(select/flatSelect/process + 超时侧输出)。
三条实操建议:
- 新作业默认 process:PatternProcessFunction 一个类收口匹配与超时,后续加侧输出/时间戳不用改结构;
- 跨事件条件用 IterativeCondition 而不是外部状态 hack:逐笔递增、比上一笔大这类规则,迭代条件是原生姿势,外部状态要自己管生命周期,checkpoint 还容易漏;
- 超时出口必须配:handleTimeout 或 timeoutFn 不写,超时的部分匹配静默丢弃——对"下单未支付"这类业务就是数据丢失。