ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

PyFlink DataStream 有状态流处理实战:ValueState 状态访问与 State TTL 配置全解析

PyFlink DataStream 有状态流处理实战:ValueState 状态访问与 State TTL 配置全解析 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文围绕 Apache Flink 官方 PyFlink 示例中的状态访问范例flink-python/docs/examples/datastream/state.rst展开该文档通过literalinclude指令直接嵌入并展示可运行示例 state_access.py 的完整源码。阅读本文后你将掌握在 PyFlink DataStream API 中编写有状态处理函数的完整流程如何使用KeyedProcessFunction配合RuntimeContext读写ValueState如何通过StateTtlConfig为状态设置生存时间TTL并选择更新、可见性与清理策略以及 PyFlink 提供的全部五种状态类型及其描述符。全文以示例代码为骨架并给出 pyflink/datastream/state.py 与 pyflink/datastream/functions.py 中的源码级佐证可直接复制运行、对照验证。一、示例定位一份以代码为主体的状态访问文档在 PyFlink 官方文档的示例集中每个主题对应一个 RST 文件与一个同名 Python 脚本例如本主题的状态访问对应文档页state.rst源码脚本state_access.pystate.rst的正文只有一条literalinclude指令把state_access.py的源码原样嵌入文档。因此在阅读本文时应把示例脚本视为文档的主体内容。脚本展示的是一个经典的按键累计求和场景以用户名为 key持续累加每个用户上报的金额并利用ValueState保存中间累计值同时为状态配置 1 秒的 TTL 生命周期。这是一个自包含、可独立运行的完整作业也是理解 PyFlink 有状态算子的最小实践范本。二、完整示例代码先给出 state_access.py 的完整源码Apache License 2.0 许可后续各小节将逐段拆解from pyflink.common import Time from pyflink.common.typeinfo import Types from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.functions import KeyedProcessFunction, RuntimeContext from pyflink.datastream.state import ValueStateDescriptor, StateTtlConfig class Sum(KeyedProcessFunction): def __init__(self): self.state None def open(self, runtime_context: RuntimeContext): state_descriptor ValueStateDescriptor(state, Types.FLOAT()) state_ttl_config StateTtlConfig \ .new_builder(Time.seconds(1)) \ .set_update_type(StateTtlConfig.UpdateType.OnReadAndWrite) \ .disable_cleanup_in_background() \ .build() state_descriptor.enable_time_to_live(state_ttl_config) self.state runtime_context.get_state(state_descriptor) def process_element(self, value, ctx: KeyedProcessFunction.Context): # retrieve the current count current self.state.value() if current is None: current 0 # update the states count current value[1] self.state.update(current) yield value[0], current def state_access_demo(): env StreamExecutionEnvironment.get_execution_environment() ds env.from_collection( collection[ (Alice, 110.1), (Bob, 30.2), (Alice, 20.0), (Bob, 53.1), (Alice, 13.1), (Bob, 3.1), (Bob, 16.1), (Alice, 20.1) ], type_infoTypes.TUPLE([Types.STRING(), Types.FLOAT()])) # apply the process function onto a keyed stream ds.key_by(lambda value: value[0]) \ .process(Sum()) \ .print() # submit for execution env.execute() if __name__ __main__: state_access_demo()代码逻辑十分直观输入是(用户名, 金额)二元组流Sum算子为每个用户维护一个 Float 类型的累计值状态每来一条数据就把金额累加到该用户的状态中并输出(用户名, 当前累计值)。三、逐步拆解一个完整的有状态作业是如何构建的3.1 创建执行环境与内存数据源env StreamExecutionEnvironment.get_execution_environment() ds env.from_collection( collection[ (Alice, 110.1), (Bob, 30.2), (Alice, 20.0), (Bob, 53.1), (Alice, 13.1), (Bob, 3.1), (Bob, 16.1), (Alice, 20.1) ], type_infoTypes.TUPLE([Types.STRING(), Types.FLOAT()]))StreamExecutionEnvironment.get_execution_environment()创建 PyFlink DataStream 执行环境。from_collection从一个 Python 列表构造数据源并通过type_info显式声明元素类型为(STRING, FLOAT)二元组——类型信息TypeInformation是后续创建状态描述符的基础PyFlink 的序列化与类型推断都依赖它。这里刻意选取了 Alice 与 Bob 两组交错数据目的是展示状态如何按 key 隔离、分别累计。3.2 key_by 分流与 process 应用ds.key_by(lambda value: value[0]) \ .process(Sum()) \ .print() env.execute()key_by(lambda value: value[0])按元组第一个元素用户名对流进行分区生成KeyedStream。注意PyFlink 的状态访问只对 KeyedStream 上的算子开放——这正是KeyedStateStore接口functions.py所强调的key/value state is only accessible if the function is executed on a KeyedStream。Flink 运行时保证同一 key 的所有元素路由到同一个并行子任务状态因此得以透明地分区、扩容与重分布。.process(Sum())将KeyedProcessFunction应用到 KeyedStream 上这是最灵活的低层处理算子可以访问 key、读写状态、注册定时器本示例未涉及定时器相关示例见 event_time_timer.py。.print()把输出打印到标准输出便于本地调试。env.execute()提交作业执行。注意execute()会阻塞直到作业结束from_collection产生的有界数据流处理完毕后作业即退出。3.3 open() 中初始化状态RuntimeContext 与 ValueStateDescriptorSum类继承自KeyedProcessFunctionfunctions.py生命周期方法open(runtime_context)在算子初始化时被调用一次状态创建必须放在这里def __init__(self): self.state None def open(self, runtime_context: RuntimeContext): state_descriptor ValueStateDescriptor(state, Types.FLOAT()) ... self.state runtime_context.get_state(state_descriptor)关键点有三个状态句柄不能在__init__中创建。RuntimeContext只有在open阶段才可用因此先置self.state None在open中通过runtime_context.get_state(descriptor)获取状态句柄。ValueStateDescriptor描述状态的元信息。它的构造签名是ValueStateDescriptor(name, value_type_info)见 state.py这里状态名为state值类型为Types.FLOAT()。状态名用于区分同一个算子内注册的多个状态类型信息用于序列化与反序列化。RuntimeContext不止提供状态访问。它还暴露了并行度、子任务编号、任务名、尝试次数等静态上下文信息functions.py但状态访问是它最核心的能力KeyedStateStore接口统一定义了五类状态的获取方法见第五节。3.4 process_element 中读写状态process_element是每条数据的处理入口签名与 Java 的ProcessFunction.processElement对应def process_element(self, value, ctx: KeyedProcessFunction.Context): current self.state.value() if current is None: current 0 current value[1] self.state.update(current) yield value[0], currentself.state.value()读取当前 key 的状态值。ValueState.value()的返回类型是Optional[T]首次访问时为None必须判空示例用current 0作为初始值见 state.py 的value()/update()定义。self.state.update(current)写回新的累计值。yield value[0], currentKeyedProcessFunction.process_element是生成器函数通过yield逐条输出下游元素。对整个数据集跑一遍Alice 的累计值依次为 110.1 → 130.1 → 143.2 → 163.3Bob 的累计值依次为 30.2 → 83.3 → 86.4 → 102.5。两个 key 的状态互不干扰这正是 keyed state 的语义。四、StateTtlConfig为状态设置生存时间示例最有技术含量的是状态 TTL 配置state_ttl_config StateTtlConfig \ .new_builder(Time.seconds(1)) \ .set_update_type(StateTtlConfig.UpdateType.OnReadAndWrite) \ .disable_cleanup_in_background() \ .build() state_descriptor.enable_time_to_live(state_ttl_config)StateTtlConfig是 PyFlink 中 TTL 配置的唯一入口完整定义见 state.py它通过流式 Builder 构建。enable_time_to_live定义在状态描述符基类StateDescriptor上state.py意味着五种状态类型都可以启用 TTL。下面逐个拆解配置维度。4.1 UpdateType何时刷新最后访问时间TTL 的判定依据是每个状态值的最后访问时间戳UpdateType决定这个时间戳何时被刷新从而延长或缩短状态的存活期。枚举定义在 state.py取值语义DisabledTTL 关闭状态永不过期OnCreateAndWrite状态创建时初始化时间戳且每次写操作时更新OnReadAndWrite与OnCreateAndWrite相同但读操作也会刷新时间戳Builder 也提供了语义化快捷方法update_ttl_on_create_and_write()与update_ttl_on_read_and_write()。示例选用OnReadAndWrite意味着只要某个 key 的状态被读到其 TTL 就会被续命1 秒——适合高频访问的热数据避免其过早过期。4.2 StateVisibility过期但尚未清理的值能否被读到TTL 过期与实际物理清理之间存在时间差StateVisibility控制这个窗口期内读操作的行为state.py取值语义ReturnExpiredIfNotCleanedUp若过期值尚未被清理仍将其返回给用户NeverReturnExpired绝不返回过期值过期即视为不存在两者的取舍是数据新鲜度与读取成本NeverReturnExpired保证读到的永远是非过期数据但每次读都需要检查时间戳。示例未显式设置采用 Builder 默认值NeverReturnExpired。4.3 TtlTimeCharacteristicTTL 的时间尺度TtlTimeCharacteristic决定 TTL 采用哪种时钟state.py取值语义ProcessingTime处理时间当前仅支持此选项当前 PyFlink 仅实现ProcessingTime即按算子本地的处理时钟判定过期可通过use_processing_time()显式指定。这与 Flink Java API 中 TTL 默认使用处理时间的行为一致。4.4 清理策略过期状态何时被物理移除StateTtlConfig的清理策略是另一个重要维度。Builder 提供三种显式策略cleanup_full_snapshot()在 checkpoint 全量快照时清理过期状态。不在线清理不影响读写路径性能但快照体积中仍会短暂保留过期数据。cleanup_incrementally(cleanup_size, run_cleanup_for_every_record)每次状态访问时批量检查cleanup_size个 key 并清理过期项若run_cleanup_for_every_recordTrue则每条记录都会触发一轮清理。源码注释state.py特别提醒该策略目前只对 Heap 状态后端生效对 RocksDB 无效且会增加单条记录的处理延迟若 Heap 后端使用同步快照全局迭代器会持有全部 key 的副本内存占用上升。cleanup_in_rocksdb_compact_filter(query_time_after_num_entries, periodic_compaction_timeNone)在 RocksDB 压缩过程中清理过期状态。压缩过滤器每处理query_time_after_num_entries条状态条目就查询一次当前时间戳periodic_compaction_time用于对不常访问的状态做定期压缩默认 30 天state.py。示例调用的是disable_cleanup_in_background()即关闭默认的后台清理。注意 Builder 的默认_is_cleanup_in_background Truestate.py后台清理对应增量清理默认参数(5, False)与 RocksDB 压缩过滤默认query_time_after_num_entries1000state.py。disable_cleanup_in_background()会把这些默认后台策略一并关闭但源码注释明确如果显式配置了cleanup_incrementally或cleanup_in_rocksdb_compact_filter此开关不会禁用它们state.py。4.5 Builder 默认值一览综合 Builder 初始化StateTtlConfig.new_builder(ttl)的默认配置为配置项默认值update_typeOnCreateAndWritestate_visibilityNeverReturnExpiredttl_time_characteristicProcessingTimeis_cleanup_in_backgroundTrue增量清理 5 条/次、RocksDB 压缩过滤器 1000 条/次显式清理策略无示例中1 秒 TTL 读改写双刷新 关闭后台清理的组合意味着状态在创建或读写后 1 秒即过期过期后读取视为不存在由于关闭了后台清理过期条目依赖后续访问或快照/压缩过程逐步移除适合演示 TTL 的语义而不过度纠结清理时机。五、PyFlink 支持的完整状态类型体系示例只用了ValueState但 PyFlink 在KeyedStateStore接口functions.py中统一定义了五种可用的 keyed state 类型对应的状态类与描述符均位于 state.py状态类型获取方法描述符典型语义ValueState[T]get_stateValueStateDescriptor(name, value_type_info)单个值接口为value()/update()/clear()ListState[T]get_list_stateListStateDescriptor(name, elem_type_info)可追加、可整体覆盖的列表接口为add()/update(list)/get()MapState[K, V]get_map_stateMapStateDescriptor(name, key_type_info, value_type_info)key-value 映射接口为get(key)/put(key, value)/remove(key)等并支持__setitem__/__delitem__语法糖ReducingState[T]get_reducing_stateReducingStateDescriptor(name, reduce_function, type_info)用ReduceFunction增量聚合状态内部只保存聚合结果AggregatingState[IN, OUT]get_aggregating_stateAggregatingStateDescriptor(name, agg_function, state_type_info)用AggregateFunction聚合允许累加器类型与输出类型不同描述符统一继承自StateDescriptor基类因此enable_time_to_live()对上述所有状态类型都可用state.py。一个函数可以注册多个不同名字的状态例如同时用ValueState保存上次事件时间、用MapState维护滑动窗口内的明细在open中分别get_state/get_map_state即可。六、运行与验证该示例不依赖外部组件在已安装 PyFlink 的 Python 环境中直接执行即可python flink-python/pyflink/examples/datastream/state_access.py本地运行时会以单进程模拟完整集群执行print()输出会打印到控制台形如(Alice, 110.1)、(Alice, 130.1)等累计结果env.execute()结束后进程正常退出。若要修改验证行为可调整collection中的数据、Time.seconds(1)的 TTL 时长或把set_update_type换成OnCreateAndWrite、去掉disable_cleanup_in_background()对比观察过期语义的变化。同目录下的 word_count.py、streaming_word_count.py、basic_operations.py 等示例覆盖了更多 DataStream 操作可对照学习。七、小结通过这个官方示例可以看到 PyFlink 有状态编程的完整链路key_by得到 KeyedStream → 自定义KeyedProcessFunction→ 在open()中用描述符 RuntimeContext.get_state()创建状态 → 在process_element()中value()/update()读写 → 用StateTtlConfig精确控制状态的生存时间与清理行为。状态描述符的创建、TTL 的启用以及五种状态类型的选择是编写正确、可容错、可扩展的有状态 PyFlink 作业的三个基本功示例脚本 state_access.py 与状态实现 state.py 均可直接作为你后续开发的参考模板。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Flink DataStream 状态编程实战Keyed State、Operator State 与状态 TTL 全解析Apache Flink DataStream 状态编程实战Keyed State、Operator State 与状态 TTL 全解析 本文是 Apache大数据流处理批处理数据工程Flink 有状态流处理编程指南Keyed State、算子状态与状态 TTL 完整实战Flink 有状态流处理编程指南Keyed State、算子状态与状态 TTL 完整实战 本指南基于当前仓库中的 Working with State 官方文大数据流处理批处理数据工程PyFlink DataStream 状态编程完全指南State、StateTtlConfig 与 StateBackend 深度解析PyFlink DataStream 状态编程完全指南State、StateTtlConfig 与 StateBackend 深度解析 PyFlink 在 p大数据流处理批处理数据工程上一篇【亲测免费】 稳定的时间戳为Whisper增稳下一篇如何快速安装SpotX-Bash面向新手的完整教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表