摘要:ProcessFunction 是 Flink 表达力最强的函数类,但"会用"和"懂原理"是两回事。这篇文章三层递进:先梳理完整 API 谱系(单流/双流/窗口/Broadcast 四个变体);再拆内部执行链路——一条数据从网络进来到回调 processElement 经过哪些环节、watermark 如何驱动 onTimer、定时器在 Heap/RocksDB 里怎么存;最后给出四个可直接运行的代码案例(会话切分、多级侧输出、动态规则、watermark 驱动实验)。读完你能解释"为什么定时器能跨重启恢复"“为什么 OutputTag 必须匿名类”“非 keyed 流为什么不能注册定时器”。
关键词:Flink ProcessFunction、KeyedProcessFunction、CoProcessFunction、BroadcastProcessFunction、定时器、TimerService、侧输出、OutputTag、内部执行链路、watermark、InternalTimerService、代码实现
一、从"会用"到"懂原理"
前两篇我们讲了函数类全景和富函数纵深。ProcessFunction 是这条线的终点——它继承 AbstractRichFunction,在富函数全部能力之上,多了定时器、侧输出、事件时间访问三个维度,是 Flink 表达力最强的 API。
但线上很多问题恰恰出在"会用不会原理":为什么定时器能跨重启恢复?为什么非 keyed 流注册定时器直接报错?为什么侧输出的 OutputTag 必须带花括号?这些问题的答案不在 API 文档里,而在 ProcessFunction 的内部执行链路中。这篇文章先给谱系,再拆链路,最后上代码。
二、API 谱系:先分清 keyed 与否、单流还是双流
ProcessFunction 家族不是一个大类,而是按"是否 keyed、单流还是双流"划分的一组 API:
| 类 | keyed | 双流 | 定时器 | 典型场景 |
|---|---|---|---|---|
| ProcessFunction | 否 | 否 | ❌ 注册即报错 | 全量流分流/告警 |
| KeyedProcessFunction | 是 | 否 | ✅ | 超时检测/会话切分(主力) |
| CoProcessFunction | 视 keyBy | 是 | keyed 后可用 | 双流对账 |
| BroadcastProcessFunction | 否 | 是(广播) | ❌ | 动态规则下发 |
| KeyedBroadcastProcessFunction | 是 | 是(广播) | ✅ | 规则+keyed 状态 |
| ProcessWindowFunction / ProcessAllWindowFunction | 视窗口 | 否 | 窗口上下文 | 窗口全量处理 |
两个容易搞混的点:
- ProcessFunction(非 keyed)没有定时器。调用
ctx.timerService().registerEventTimeTimer(...)会直接抛UnsupportedOperationException(“Registering timers is only supported on a keyed streams”)。因为定时器必须挂在 key 上才能跨 checkpoint 恢复、按 key 分发——没有 key,就没有定时器的存储和路由基础。需要定时器,先keyBy。 - BroadcastProcessFunction 的广播状态只有一边可写。
processBroadcastElement里可以读写广播状态,processElement(事件侧)只能读——规则流负责更新,事件流只消费。这个单向约束是刻意的:广播状态在每条并行实例上都有一份副本,只允许规则侧变更才能保持一致。
三、内部执行链路:一条记录和一个 watermark 的两条路径
写stream.keyBy(...).process(new KeyedProcessFunction<...>() {...})时,框架做了比你想的更多的事。
3.1 数据路径:processElement 被谁调用
你写的 KeyedProcessFunction 并不会直接被流处理引擎调用。真实链路是:
- 网络输入反序列化成
StreamRecord(value + timestamp); - 任务运行循环(
StreamInputProcessor)逐条把记录喂给算子链; - 内部算子
KeyedProcessOperator.processElement(record)收到记录——它才是真正持有你函数类的对象; - 关键一步:
setCurrentKey(record.getKey()),让KeyedStateBackend把"当前 key"切到这条记录的 key; - 之后才调用你的
processElement(value, ctx, out)。
ctx不是接口魔法,它是算子内部类ContextImpl的实例:ctx.timestamp()读的就是 StreamRecord 上的时间戳,ctx.getCurrentKey()读的是 KeyedStateBackend 的当前 key,ctx.timerService()拿到的是算子持有的定时器服务。
第 4 步是整条链路最容易被忽略、也最重要的设计:key 上下文 = 状态正确性的前提。你在 processElement 里读写 ValueState,之所以操作的是"当前这条记录所属 key"的状态,就是因为引擎在回调前做了 setCurrentKey。onTimer 回调同理——触发定时器前,引擎会setCurrentKey(timer.key),所以定时器回调里读状态,读到的是定时器所属 key的状态,而不是什么"当前数据"的。
3.2 定时器路径:watermark 如何驱动 onTimer
事件时间定时器不是"到点自动响",而是由 watermark 推动:
- watermark 事件到达任务运行循环;
InternalTimeServiceManager.advanceWatermark(wm)逐算子推进水位;- 每个算子的
InternalTimerServiceImpl从定时器队列里弹出所有 timestamp ≤ wm 的定时器(队列按 key-group 分桶、桶内按时间有序); - 对每个弹出的定时器:
setCurrentKey(timer.key)→ 调用你的onTimer(ts, ctx, out)。
所以 onTimer 触发的准确语义是:watermark 越过了这个时间点。而处理时间定时器走另一条线(ProcessingTimeService 按本地时钟轮询),到点即触发,不管 watermark——这也是为什么重启后处理时间定时器"不追补"、事件时间定时器"原样恢复"。
四、定时器与侧输出:两个"状态级"能力的内部机制
4.1 定时器:持久化的闹钟
- 注册:
registerEventTimeTimer(ts)内部封装成InternalTimer(key, namespace, ts)进入定时器队列; - 去重:同 key + 同时间戳只保留一个(namespace 用于区分窗口等场景);
- 存储:Heap StateBackend 用
FlinkPriorityQueue(堆内有序队列);RocksDB 用KeyGroupedInternalPriorityQueue——定时器也落盘,量大时和状态一起占磁盘; - 恢复:定时器随 checkpoint 序列化,作业重启后原样恢复,事件时间语义不丢。
由此推出两个工程结论:定时器数量 = key 数 × 时间点,超量会拖垮 checkpoint 和恢复(用时间对齐/单调推进收敛);删除定时器必须 key + namespace + ts 与注册时完全一致,否则删不掉、只能靠状态判空兜底。
4.2 侧输出:带标签的旁路管道
- 定义:
new OutputTag<X>("name") {}——必须匿名类带花括号。花括号让 Java 保留泛型类型信息(TypeInformation),侧输出反序列化需要它;不带{}的new OutputTag<>("name")类型擦除后拿不到类型,运行期报错; - 发射:
ctx.output(tag, value)内部走SideOutputDataOutput旁路通道; - 输出:旁路记录随主流一起输出、一起参与 checkpoint 与水位线传递,不丢数据;
- 取流:下游
dataStream.getSideOutput(tag)取出,tag 需与发射时同一个。
五、代码实现:四个可直接运行的案例
5.1 ProcessFunction 基础(非 keyed):分流 + 定时器报错现场
// 场景:设备日志全量流,按级别分流 + 数每条日志的处理时间DataStream<String>main=logs.process(newProcessFunction<Log,String>(){// 侧输出标签(static 匿名类形式,保留类型信息)privatestaticfinalOutputTag<String>ERROR_TAG=newOutputTag<String>("error"){};privatestaticfinalOutputTag<String>WARN_TAG=newOutputTag<String>("warn"){};@OverridepublicvoidprocessElement(Loglog,Contextctx,Collector<String>out){if(log.level>=40){ctx.output(ERROR_TAG,log.msg);// 错误 → 侧输出}elseif(log.level>=30){ctx.output(WARN_TAG,log.msg);// 警告 → 另一个侧输出}else{out.collect(log.msg);// 正常 → 主流}// ctx.timestamp():事件时间模式下非 null,处理时间模式下为 nulllongts=ctx.timestamp()==null?-1:ctx.timestamp();// 想注册定时器?这里会抛 UnsupportedOperationException:// ctx.timerService().registerEventTimeTimer(ts);}});// 下游分别取三条流DataStream<String>errorStream=main.getSideOutput(ERROR_TAG);DataStream<String>warnStream=main.getSideOutput(WARN_TAG);注意OutputTag定义成static final字段——和富函数篇的序列化教训一脉相承:非静态匿名类会捕获外部 this,OutputTag 定义在算子内部更要注意 static。
5.2 KeyedProcessFunction 完整案例:用户会话切分
场景:用户行为流,10 分钟无操作视为会话结束,输出会话的行为数和时长。这是定时器+状态+侧输出的全家桶案例(与订单超时不同,这里每次事件都要重置定时器——会话 gap 是滑动重置的):
// keyBy(userId) 之后publicclassSessionSplitterextendsKeyedProcessFunction<Long,UserAction,Session>{privateValueState<Long>firstTs;// 会话开始时间privateValueState<Integer>count;// 会话内行为数privatestaticfinallongGAP=10*60*1000L;privatestaticfinalOutputTag<UserAction>LATE_TAG=newOutputTag<UserAction>("late"){};@Overridepublicvoidopen(Configurationparameters){firstTs=getRuntimeContext().getState(newValueStateDescriptor<>("first",Long.class));count=getRuntimeContext().getState(newValueStateDescriptor<>("cnt",Integer.class));}@OverridepublicvoidprocessElement(UserActiona,Contextctx,Collector<Session>out)throwsException{Longfirst=firstTs.value();if(first==null){// 会话开始:记录起点,注册 GAP 后的定时器firstTs.update(a.ts);count.update(1);ctx.timerService().registerEventTimeTimer(a.ts+GAP);}else{count.update(count.value()+1);// 新行为刷新会话:删掉旧定时器,重新注册(gap 滑动重置)ctx.timerService().deleteEventTimeTimer(first+GAP);ctx.timerService().registerEventTimeTimer(a.ts+GAP);firstTs.update(a.ts);}}@OverridepublicvoidonTimer(longts,OnTimerContextctx,Collector<Session>out)throwsException{// 距最后一次行为已过 GAP,会话结束Integerc=count.value();if(c!=null){out.collect(newSession(ctx.getCurrentKey(),c,ts-GAP,ts));firstTs.clear();count.clear();}}}这个案例的工程细节:每次行为都"删旧 + 注册新",代价是定时器注册/删除频繁,但语义正确(gap 从最后一次行为算起);onTimer里ts就是定时器触发时间,ts - GAP即会话起点——不需要额外状态记会话结束时间。
5.3 BroadcastProcessFunction:动态规则下发
场景:风控事件流 + 规则流(阈值实时更新),规则变化不重启作业:
// 规则流 broadcast 到所有并行实例MapStateDescriptor<String,Rule>RULE_DESC=newMapStateDescriptor<>("rules",String.class,Rule.class);DataStream<Rule>ruleStream=env.fromElements(newRule("amount",10000));BroadcastStream<Rule>bcRules=ruleStream.broadcast(RULE_DESC);events.connect(bcRules).process(newBroadcastProcessFunction<Event,Rule,Alert>(){@OverridepublicvoidprocessElement(Evente,ReadOnlyContextctx,Collector<Alert>out){// 事件侧:只读广播状态(ReadOnlyContext 只暴露只读接口)Rulerule=ctx.getBroadcastState(RULE_DESC).get("amount");if(rule!=null&&e.amount>rule.threshold){out.collect(newAlert(e,rule));}}@OverridepublicvoidprocessBroadcastElement(Ruler,Contextctx,Collector<Alert>out){// 规则侧:可写广播状态——热更新阈值ctx.getBroadcastState(RULE_DESC).put("amount",r);}});两个设计要点:事件侧上下文是ReadOnlyContext——编译期就禁止你写广播状态,这是 API 层面做的正确性约束;规则流要低吞吐(一条规则广播到所有实例成本不低),别把高频数据流塞进 broadcast。
5.4 底层验证实验:打印 watermark 与 onTimer 的触发关系
理解第三节链路的最好方式,是做一个观察实验——在 onTimer 里打印当前 watermark 与触发时间:
// 事件时间 + 每 2 秒一个 watermark,观察触发条件env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);// 1.13 前写法DataStream<String>out=stream.assignTimestampsAndWatermarks(WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(2)).withTimestampAssigner((s,ts)->parseTs(s))).keyBy(s->s.split(",")[0]).process(newKeyedProcessFunction<String,String,String>(){@OverridepublicvoidprocessElement(Stringv,Contextctx,Collector<String>out){longt=ctx.timestamp();// 注册 t+5000 的事件时间定时器(同一 key 同一 ts 只注册一次)ctx.timerService().registerEventTimeTimer(t+5000);out.collect("registered timer at "+t);}@OverridepublicvoidonTimer(longts,OnTimerContextctx,Collector<String>out){// 触发时 watermark 一定 >= ts;打印两者关系out.collect("onTimer ts="+ts+" watermark="+ctx.timerService().currentWatermark());}});// 观察日志:onTimer 触发的时间点 = watermark 越过 ts 的那一次推进观察结论会非常直观:onTimer 的触发时间不是"注册后 5 秒",而是"watermark 推进到 ts 之后"——如果 watermark 卡住(数据停了、乱序超界),定时器永远不触发。这也是为什么事件时间定时器不适合"物理超时兜底",那种场景要用处理时间定时器或外部兜底任务。
六、实战避坑清单
- 非 keyed 流注册定时器 → UnsupportedOperationException。要定时器先 keyBy;ProcessFunction 只做分流/告警。
- OutputTag 必须匿名类
{},且定义成 static final——既保类型信息,又避免捕获外部 this。 - 删除定时器必须完全一致:key + namespace + ts 三个要素和注册时相同,否则删不掉。
- 定时器数量要管理:= key 数 × 时间点,RocksDB 下直接吃磁盘;用时间对齐(注册到窗口边界)或单调推进(只保留最新)收敛。
- onTimer 里读状态读的是定时器所属 key:这是引擎 setCurrentKey 的设计,也是多 key 场景写对逻辑的前提。
- watermark 卡住 = 事件时间定时器不触发:别用事件时间定时器做物理超时;要兜底用处理时间或外部调度。
- BroadcastProcessFunction 事件侧只读:想改广播状态必须在 processBroadcastElement 里;规则流别高频。
七、总结:我的判断
ProcessFunction 的"底层"其实就三件事:key 上下文切换(setCurrentKey)、定时器队列(InternalTimerService)、侧输出通道(SideOutputDataOutput)——它们都是算子内部实现的支撑,你的函数只是被回调的出口。理解这三件事,就能解释这个 API 家族几乎所有的行为特性:
- 定时器能跨重启恢复 → 因为它是状态(落 StateBackend);
- 非 keyed 不能注册定时器 → 因为定时器要挂 key 才能存储和路由;
- 定时器回调里能正确读状态 → 因为引擎先 setCurrentKey(timer.key)。
给三条实操建议:主力用 KeyedProcessFunction(80% 的状态+定时场景它都能覆盖);侧输出优先于 filter 二次遍历(一条流分叉更清晰);定时器先做数量评估再上线(写之前算一下 key 数 × 时间点,超量先对齐)。