
摘要上一篇讲函数类全景时RichFunction 只是第二家族这篇把它单独挖到底。文章从四个维度拆解富函数接口与继承结构、从实例化到 close 的完整生命周期时序状态恢复先于 open、链内上游先 open、RuntimeContext 六类能力状态/指标/分布式缓存/算子状态/配置/累加器以及一组真实的序列化与执行陷阱——含非静态匿名类隐式捕获外部 this 导致序列化爆炸、异步 I/O 必须用连接池等。看完能回答为什么构造器拿不到 RuntimeContext“open 里到底能做什么不能做什么”operator state 和 keyed state 有什么区别这类进阶问题。关键词Flink 富函数、RichFunction、RichMapFunction、生命周期、open/close、RuntimeContext、分布式缓存、operator state、CheckpointedFunction、RichAsyncFunction、序列化陷阱、异步IO一、从会用到懂执行时机大部分 Flink 开发者都会写RichMapFunctionopen 里建连接池map 里干活close 里释放。但下面这几个问题能答上来的就少了一半为什么连接池必须放在 open() 里建构造器里建不行open() 被调用时状态恢复完了吗从 checkpoint 重启后open 里看到的是空状态还是恢复好的状态算子链上有多个 Rich 函数open 的执行顺序是什么new RichMapFunction() {...}写在类的实例方法里为什么作业提交时突然报序列化异常这些问题的答案都藏在富函数的执行时序和上下文注入机制里。上一篇我们看的是函数类全景四大家族的划分这篇下沉到 RichFunction 内部把生命周期、RuntimeContext、序列化机制一次讲透。二、富函数家族结构接口四方法 模板基类 算子变体富函数不是一个类而是一族。结构分三层第一层RichFunction 接口org.apache.flink.api.common.functions。定义四个方法open(Configuration)、close()、setRuntimeContext(RuntimeContext)、getRuntimeContext()。前两个是生命周期钩子由框架回调后两个是上下文注入与读取注意setRuntimeContext 是框架内部调的你只调 getRuntimeContext。第二层AbstractRichFunction 模板基类。持有runtimeContext字段open/close 给空实现——所有算子的 Rich 变体都继承它你写的时候只需要覆盖 open/close 里需要的那部分。第三层各算子的 Rich 变体。RichMapFunction/RichFlatMapFunction/RichFilterFunction/RichCoFlatMapFunctionconnect 双流/RichSinkFunction/RichSourceFunction/RichAsyncFunction异步 I/O 专用。变体只改处理方法的签名map/flatMap/filter生命周期三件套完全一致——所以学会一个全都会了。还有一个容易漏掉的组合CheckpointedFunction 接口。它和 RichFunction 是正交的可以同时实现RichFunction 管生命周期和上下文CheckpointedFunction 管非 keyed 的算子状态initializeState/snapshotState两个方法。记一个口诀keyed state 用getRuntimeContext().getState()operator state 用 CheckpointedFunction。三、生命周期时序状态恢复先于 open这是最容易被忽略的细节单个并行子任务的函数实例完整时序是七步实例化Task 启动算子链内的函数反序列化构造器执行。此刻 RuntimeContext还不存在——构造器里调getRuntimeContext()直接 NPEsetRuntimeContext 注入框架内部完成发生在 open 之前initializeState 状态恢复keyed state、operator state、定时器从 StateBackend 恢复先于 open 完成open(Configuration)每个并行子任务各执行一次算子链内按上游 → 下游顺序逐个 openprocessElement 循环单线程串行处理checkpoint 周期快照贯穿处理期状态持久化函数实例本身不参与快照close()任务正常结束时释放资源故障/强制 cancel 不保证执行。用一段带日志的代码验证这个时序生产排查时这也是标准手法// 验证构造 → 注入 → 恢复 → open → 处理 → close 的顺序DataStreamStringoutsource.flatMap(newRichFlatMapFunctionString,String(){Overridepublicvoidopen(Configurationparameters)throwsException{// 此时状态已恢复从 checkpoint 重启这里读到的就是恢复好的值Longrestoredstate.value();// state: ValueStateLongLOG.info(open: restored{}, subtask{},restored,getRuntimeContext().getIndexOfThisSubtask());// 能拿到并行度、任务名、全局参数、缓存文件——上下文已完整注入}OverridepublicvoidflatMap(Stringvalue,CollectorStringout){// 每条记录调用单线程无需加锁state.update(System.currentTimeMillis());out.collect(value);}Overridepublicvoidclose()throwsException{LOG.info(close: subtask{},getRuntimeContext().getIndexOfThisSubtask());}});对应到生产实践就是三个一定连接池、HTTP client、加载的缓存文件一定放 open()。它们是每个子任务一份的正好和 open 的执行粒度对齐状态句柄一定在 open() 里注册一次、存成员变量。每次处理都调getState()会重复创建句柄对象纯开销一定不要用构造器做初始化。除了拿不到上下文还有个隐蔽问题函数实例要序列化分发构造器里建好的连接在反序列化后已经无效连接绑定的是原进程的资源。四、RuntimeContext 六类能力状态、指标、缓存、算子状态、配置getRuntimeContext()返回的接口能力可以分成六类前四类生产最常用4.1 keyed state 句柄最常见ValueState / ListState / MapState / ReducingState / AggregatingState按 key 隔离落 StateBackend随 checkpoint 持久化。前提是 keyBy 之后——非 keyed 流上用会直接报错。4.2 指标 MetricGroup可观测性// open() 里注册processElement 里打点——比日志统计可靠得多CountercntgetRuntimeContext().getMetricGroup().counter(hit_count);// 计数累计值MetermetergetRuntimeContext().getMetricGroup().meter(hit_rate,newMeterView(cnt));// 速率每秒增量GaugeLonggaugegetRuntimeContext().getMetricGroup().gauge(pending,()-pendingSize);// 瞬时值lambda 读取指标会聚合到 JobManager可接 Prometheus/Grafana 出面板。计数类指标优先用 MetricGroup.counter别用累加器 Accumulator——累加器与 checkpoint 快照存在兼容性问题官方也推荐 metric 优先累加器只留给作业结束时拿结果的少数场景。4.3 分布式缓存 DistributedCache词表/规则/模型文件// 驱动端注册提交前env.registerCachedFile(hdfs://namenode:8020/data/blacklist.txt,blacklist);// 函数端 open() 里读取FilefilegetRuntimeContext().getDistributedCache().getFile(blacklist);文件会被分发到所有 TaskManager 的本地磁盘每个子任务本地一份open 里加载进内存即可——避免了每条记录远程读的灾难。注意两点文件更新要重启作业才生效超大文件几百 MB 以上别走缓存用外部存储 惰性加载。4.4 operator state非 keyed 状态与 CheckpointedFunction这是很多文章一笔带过、但生产极有用的能力。非 keyed 流比如 Source、无 keyBy 的算子没有 key 概念但可以持有按算子实例分片的算子状态配合 CheckpointedFunction 实现跨并行度可恢复的计数/位点// 示例全局事件计数非 keyed 流也能做、能随 checkpoint 恢复publicclassGlobalCounterextendsRichFlatMapFunctionString,LongimplementsCheckpointedFunction{privatelonglocalCount0;// 普通成员本实例的计数privateListStateLongcheckpointed;// 算子状态checkpoint 用OverridepublicvoidinitializeState(FunctionInitializationContextctx)throwsException{// 恢复/初始化算子状态isRestored 区分首次启动还是从 checkpoint 恢复checkpointedctx.getOperatorStateStore().getListState(newListStateDescriptor(count,Long.class));if(ctx.isRestored()){for(Longv:checkpointed.get())localCountv;}}OverridepublicvoidsnapshotState(FunctionSnapshotContextctx)throwsException{checkpointed.clear();checkpointed.add(localCount);// 快照时把当前值写入算子状态}OverridepublicvoidflatMap(Stringvalue,CollectorLongout){localCount;out.collect(localCount);}}和 keyed state 的差异要记清楚keyed state按 key 分片key 分布到哪个子任务状态就在哪operator state按算子实例分片每个并行子任务一份重启时可选择 list 均分或 union 合并再分配另外operator state 不支持 TTL需要自己管理清理。典型应用Kafka 自定义 Source 的 offset 位点记录、跨并行度的全局计数。4.5 配置与任务信息getRuntimeContext().getExecutionConfig().getGlobalJobParameters()读提交时注入的全局参数改参数需重启getNumberOfParallelSubtasks()/getIndexOfThisSubtask()拿到分片信息可以做按子任务编号分目录落盘这类操作getTaskName()/getAttemptNumber()用于日志里定位。4.6 累加器历史遗留getAccumulator(name)/addAccumulator()作业结束后结果聚合回驱动端。前面说过计数优先用 metric别拿累加器当状态用——它既不参与 checkpoint语义上也不是为跨记录累积设计的。五、序列化陷阱为什么作业提交时突然炸了富函数实例要跨 JobManager → TaskManager 序列化分发这带来一类高频坑都围绕什么被序列化了。陷阱 1非静态匿名类隐式捕获外部 this 最常见的翻车现场在某个类的实例方法里写publicclassOrderProcessor{privateRedisClientredisClient;// 不可序列化的成员publicvoidrun(DataStreamStrings){s.map(newRichMapFunctionString,String(){// ⚠️ 匿名类OverridepublicStringmap(Stringv){returnenrich(v);}});}}匿名类隐式持有外部实例OrderProcessor.this的引用。作业提交时 Flink 序列化这个匿名类会把OrderProcessor整个实例一起序列化——redisClient不可序列化直接NotSerializableException就算外部类恰好可序列化也会把一大棵对象树复制到每个 TaskManager静默拖慢提交甚至 OOM。修复要么把匿名类改成static 修饰的具名内部类static 内部类不持有外部 this要么把函数类独立成文件要么把外部类实现 Serializable 并确保成员都可序列化。写富函数前先问一句这个类是不是 static / 独立类陷阱 2transient 资源三件套连接池、线程池、HTTP client 一律private transient open 里建 close 里关。漏标 transient 会在提交时炸标了但没在 open 里重建运行期 NPE——两个错误成对出现检查顺序先看声明再看 open。陷阱 3lambda 捕获外部非序列化变量map(x - doSth(externalClient, x))里捕获了externalClientlambda 序列化时同样要序列化它。规则和陷阱 1 一致函数内引用的外部对象全都会被拖进序列化。六、RichAsyncFunction异步 I/O 的线程模型与连接池异步 I/O 是富函数家族里特殊的一员——它是唯一回调并发执行的富函数线程模型和其他成员完全不同坑也最典型publicclassAsyncRedisLookupextendsRichAsyncFunctionString,Enriched{privatetransientJedisPoolpool;// ⚠️ 必须连接池asyncInvoke 并发调用Overridepublicvoidopen(Configurationparameters){poolnewJedisPool(newJedisPoolConfig(),redis-1,6379);// 初始化连接池发生在 open主线程asyncInvoke 在 I/O 线程池并发执行}OverridepublicvoidasyncInvoke(Stringkey,ResultFutureEnrichedresultFuture){try{try(Jedisjedispool.getResource()){// 每次从池里借用完归还resultFuture.complete(Collections.singleton(enrich(key,jedis)));}}catch(Exceptione){resultFuture.completeExceptionally(e);// 失败要显式回调否则任务悬挂}}Overridepublicvoidclose()throwsException{if(pool!null)pool.close();}}三个要点连接必须从连接池获取——asyncInvoke 被多个线程并发调用单个 Jedis 实例非线程安全失败必须 completeExceptionally——不回调 ResultFuture这一批数据会一直悬挂最终触发超时/背压open 里初始化一次close 里释放——生命周期三件套对异步函数同样成立。异步 I/O 的本质收益是单条数据的等待时间被并发请求覆盖吞吐提升一个量级但代价是顺序性丢失——需要保序的场景要设置AsyncDataStream.OutputMode.ORDERED。七、选型阶梯什么时候从富函数升级到 ProcessFunction富函数不是终点它和普通函数、ProcessFunction 构成三级能力阶梯第 1 级 普通函数纯转换/过滤零样板可链化。不能生命周期、状态、指标、缓存第 2 级 富函数 open/close 生命周期、 RuntimeContext 全能力。覆盖线上 80% 的算子——连接池、词表、指标、状态读写Rich 版本都够第 3 级 ProcessFunction富函数能力之上再 定时器、侧输出、事件时间访问。注意 ProcessFunction 本身继承自 AbstractRichFunction——它是富函数的超集不是另一个东西。升级触发点只有三个要等 watermark 做超时判断定时器、要把迟到/异常数据旁路侧输出、要按 key 定时触发。没有这三个需求Rich 函数就是性价比最高的选择——ProcessFunction 的样板代码和心智负担并不低滥用会让代码可读性下降。八、总结我的判断富函数的本质是 Flink 把初始化和处理分离这个工程原则做成了框架级能力open 是资源挂载点close 是资源回收点RuntimeContext 是唯一的外部能力出口。理解它的关键是执行时序——状态恢复先于 open、链内上游先 open、每个子任务独立一份这三条解释了这个家族 80% 的行为。给三条实操建议默认用 Rich 版本线上 map/flatMap/sink 一律考虑 Rich 变体哪怕暂时用不到 open——需求来了不用改结构序列化意识前置函数类写成 static 或独立类资源字段 transient检查清单在写代码时过一遍而不是等提交报错按需升级到 ProcessFunction只有定时器/侧输出/事件时间三个需求值得升级别为了显得高级把简单的转换写成 ProcessFunction。