PyFlink DataStream API 完全指南:从基础流转换到窗口、连接与广播流
发布时间:2026/9/25 2:25:13来源:尧图网络
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本文基于 Apache Flink 官方仓库中 flink-python/docs/reference/pyflink.datastream/datastream.rst 的 API 参考骨架并结合 data_stream.py 的源码实现系统讲解 PyFlink DataStream API 的九大核心类型DataStream、DataStreamSink、KeyedStream、CachedDataStream、WindowedStream、AllWindowedStream、ConnectedStreams、BroadcastStream与BroadcastConnectedStream。读完本文你将掌握 PyFlink 流式作业中算子配置、流式转换、按键分区、窗口计算、双流连接与广播状态的核心用法并理解每个 API 背后的运行机制。一、DataStream同类型元素的无限数据流DataStream是 PyFlink DataStream API 的基石。源码中将其定义为由相同类型的元素组成的流data_stream.py并可通过map、filter等转换操作派生出新的 DataStream DataStream.map(MapFunctionImpl()) DataStream.filter(FilterFunctionImpl())DataStream内部持有 Java 端org.apache.flink.streaming.api.datastream.DataStream的引用self._j_data_stream所有 Python API 最终都委托给 Java 算子执行PyFlink 因此天然具备与 Java DataStream API 对等的表达能力。1.1 算子标识与命名name / uid / set_uid_hash方法作用关键语义get_name()/name(name)获取/设置算子在可视化与运行日志中的名称用于 UI 展示与日志检索uid(uid)设置算子 ID用于跨作业提交如从 savepoint 启动时对齐算子每个转换与作业内必须唯一否则提交失败set_uid_hash(uid_hash)直接指定 JobVertexID展示于日志与 Web UI仅在默认哈希机制失效时作为 workaround 使用不能给算子链中间节点指定否则作业失败set_description(description)设置算子详细描述1.15.0 引入用于 JSON plan 与 Web UI不进入日志与指标源码 data_stream.py 中set_uid_hash的注释特别强调它是 Flink 版本升级、作业结构调整导致自动哈希变化时恢复状态与目标算子映射关系的应急手段——把旧日志里拿到的哈希通过此方法重新指定即可找回丢失的映射。1.2 并行度与资源组set_parallelism / set_max_parallelism / force_non_parallel / slot_sharing_groupset_parallelism(parallelism)设置当前算子的并行度L139-L147。set_max_parallelism(max_parallelism)设置算子最大并行度它定义动态扩缩容的上限也决定分区状态partitioned state的 key group 数量L149-L160。force_non_parallel()把并行度与最大并行度都强制为 1且后续不可再设为非 1L183-L191。适用于connect等要求单实例的操作。slot_sharing_group(slot_sharing_group)设置槽共享组。同一组内的并行实例会被尽可能放置到同一个 TaskManager 槽中未显式设置时继承输入算子的组可通过传入字符串default或 SlotSharingGroup 对象指定源码 L234-L253 支持字符串与带资源规格的对象两种形式。1.3 运行时行为调优set_buffer_timeout / start_new_chain / disable_chainingset_buffer_timeout(timeout_millis)设置数据在部分填满的缓冲区中最多滞留多久再发送到网络L193-L209。-1表示使用默认值0表示不缓冲、所有记录立即发送。较小的超时降低尾延迟但可能影响吞吐——1ms 的超时在高并行度下仍能维持高吞吐。start_new_chain()从当前算子开始一条新的算子链L211-L219。disable_chaining()关闭当前算子的链接优化L221-L232。也可用StreamExecutionEnvironment.disableOperatorChaining()关闭全作业链但源码注释提醒出于性能考虑不推荐全局关闭。1.4 数据流转换map / flat_map / filter / process四个转换算子都既接受可调用对象Python 函数也接受对应的 PyFlink 函数类MapFunction、FlatMapFunction、FilterFunction、ProcessFunctionmap(func, output_typeNone)每个元素恰好输出一个元素L273-L314。源码内部把它包装成MapProcessFunctionAdapter复用process通道实现。注意不指定output_type时输出数据会以 pickle 原始字节数组的形式序列化。flat_map(func, output_typeNone)每个元素可输出任意数量含 0 个元素L315-L355。filter(func)仅保留函数返回True的元素L415-L455。process(func, output_typeNone)最底层的转换原语每个元素可输出 0 个或多个元素并能访问上下文L641-L664。上述 map/flat_map/filter 均是它的语法糖封装。1.5 键控与分区key_by / partition_custom / 分区器家族key_by(key_selector, key_typeNone)按指定键把流分区为 KeyedStreamL357-L413。其实现颇具 PyFlink 特色先用AddKeyProcessFunction 把元素包装为Row(key, value)再调用 Java 端keyBy(JKeyByKeySelector)完成分区。不指定key_type时默认使用Types.PICKLED_BYTE_ARRAY()即 pickle 字节数组。partition_custom(partitioner, key_selector)使用自定义分区器按单字段键分区L716-L809。源码中分区函数会从任务参数NUM_PARTITIONS读取分区数非法时抛出ValueError仅支持单字段键即选择器不能返回元组。分区方法语义源码位置shuffle()输出元素随机均匀分发到下游L527-L534rescale()轮询分发到下游的子集实例子集大小取决于上下游并行度的倍数关系L555-L573rebalance()轮询分发到下游全部实例L575-L582forward()元素转发到下游的本地子任务L584-L591broadcast()元素广播到下游每个并行实例L593-L639其中broadcast有一个重载形态传入MapStateDescriptor时返回 BroadcastStream可用于后续connect map_state_desc1 MapStateDescriptor(state1, Types.INT(), Types.INT()) map_state_desc2 MapStateDescriptor(state2, Types.INT(), Types.STRING()) broadcast_stream ds1.broadcast(map_state_desc1, map_state_desc2) broadcast_connected_stream ds2.connect(broadcast_stream)1.6 时间与窗口window_all / assign_timestamps_and_watermarkswindow_all(window_assigner)对非按键分组的流开窗返回 AllWindowedStream1.16.0 引入L457-L471。assign_timestamps_and_watermarks(watermark_strategy)为元素分配时间戳并生成 watermark 以推进事件时间L666-L714。若用户指定了自定义TimestampAssigner源码会分三步走先用TimestampAssignerProcessFunctionAdapter提取时间戳并包成TUPLE([原类型, LONG])再交给 Java 端CustomTimestampAssigner分配时间戳与 watermark最后用RemoveTimestampMapFunction移除附加的时间戳字段若未指定则直接透传 Java 的WatermarkStrategy。1.7 双流合并与连接union / connectunion(*streams)将多个同类型DataStream 合并L473-L493。源码会自动处理KeyedStream参数取_values()后再合并。connect(ds)连接两个可能不同类型的流返回 ConnectedStreams若ds是 BroadcastStream 则返回 BroadcastConnectedStream1.16.0 起支持L495-L525。1.8 输出与执行add_sink / sink_to / print / execute_and_collectadd_sink(sink_func)为流添加SinkFunction型 sinkL811-L819。sink_to(sink)添加新版Sink接口的 sink如 Kafka/Pulsar 等 connectors 模块中的实现并支持SupportsPreprocessing预处理转换L821-L837。print(sink_identifierNone)把流写到标准输出L869-L884。注意它输出到运行代码的机器即 Flink worker的 stdout且不具备容错能力。内部会先调用_align_output_type把 pickle 对象或 Row 转成可读字符串。execute_and_collect(job_execution_nameNone, limitNone)触发分布式执行并通过 Flink REST API 把事件拉回当前进程L839-L867。不传limit时返回CloseableIterator必须关闭以释放集群资源传limit时返回 Python list。执行前会自动调用PythonOperatorChainingOptimizer做算子链接优化L913-L924。1.9 侧输出与缓存get_side_output / cacheget_side_output(output_tag)取出算子通过指定OutputTag发出的旁路数据流1.16.0 引入L886-L897。通常与WindowedStream.side_output_late_data配合收集迟到数据。cache()缓存变换的中间结果返回 CachedDataStream1.16.0 引入L899-L911。仅支持有界流当前仅支持 block 模式缓存在该中间结果首次被计算时惰性生成并在StreamExecutionEnvironment关闭时清除。1.10 类型与执行环境get_type()返回流的TypeInformationL162-L168所有output_type参数都可借此推导。get_execution_environment()/get_execution_config()返回创建该流的执行环境及其配置L170-L181。二、KeyedStream按键分区流KeyedStream表示算子状态按键分区通过KeySelector的 DataStreamdata_stream.py。除了shuffle、forward、key_by等分区类操作不可用外普通 DataStream 的转换都可使用且reduce、sum等归约类操作只作用于相同键的元素。2.1 键控转换与归约map/flat_map/filter/process与普通流等价但内部使用KeyedProcessFunction适配器包装可访问键控上下文L1113-L1195、L1616。reduce(func)对每个键分组内元素做归约L1197。sum(position0)、min(position0)、max(position0)、min_by(position0)、max_by(position0)按字段位置int或字段名str聚合默认取第 0 个字段L1420、L1530、L1569。min_by/max_by返回整条记录而min/max只返回对应字段值。key_by(key_selector)对已键控流重新按键如把复合键拆开再分组。union/connect/partition_custom/print/add_sink与普通流一致源码会自动适配键控流的内部结构。2.2 键控窗口与计数窗口window(window_assigner)按键开窗返回 WindowedStreamL1644。count_window(size, slide0)按元素计数开窗L1658slide0时退化为滚动计数窗口。窗口分配器定义在 window.py如TumblingEventTimeWindows、SlidingProcessingTimeWindows、TumblingCountWindows等。三、DataStreamSink流拓扑的收尾节点DataStreamSink表示流式拓扑中的 Sink 节点data_stream.py由add_sink/sink_to/print返回。只有添加了 Sink 的流才会在调用StreamExecutionEnvironment.execute()时真正被执行。它支持与 DataStream 类似的元信息配置L978-L1089方法语义name(name)设置 sink 名称用于可视化与日志uid(uid)设置算子 ID跨作业对齐如 savepoint 恢复set_uid_hash(uid_hash)直接指定 JobVertexIDset_parallelism(parallelism)设置 sink 并行度set_description(description)设置详细描述1.15.0 引入用于 JSON plan 与 Web UIdisable_chaining()关闭该算子的链接优化slot_sharing_group(group)设置槽共享组字符串或SlotSharingGroup对象四、CachedDataStream可复用缓存流CachedDataStream表示中间结果会被缓存的 DataStreamdata_stream.py缓存在该中间结果第一次被计算时生成后续使用同一CachedDataStream的作业可直接复用缓存避免重复计算。在方法集上它几乎完整继承了DataStream的能力get_type、map、flat_map、key_by、filter、window_all、union、connect、shuffle、project、rescale、rebalance、forward、broadcast、process、assign_timestamps_and_watermarks、partition_custom、add_sink、sink_to、execute_and_collect、print、get_side_output并额外提供cache()确保后续结果被缓存与DataStream.cache()语义一致。invalidate()显式使缓存失效。get_execution_environment()/set_description()获取执行环境、设置算子描述。使用方式示例在多个作业间复用中间结果cached_stream source.map(lambda x: x * 2).cache() # 后续作业可复用 cached_stream避免重新计算 map 结果五、WindowedStream按键窗口流WindowedStream表示元素按键分组、且每个键的流被WindowAssigner切成窗口的数据流data_stream.py。窗口发射由Trigger决定窗口按每个键独立评估因此不同键的窗口可以在不同时刻触发。源码特别指出L1841-L1843WindowedStream 纯属 API 层面的构造运行时它会被与 KeyedStream 和窗口算子折叠为单个算子执行——这正是窗口聚合性能的根源。5.1 窗口行为配置方法语义trigger(trigger)设置触发窗口发射的TriggerL1859-L1864。不设置时使用 WindowAssigner 的默认 Triggerallowed_lateness(time_ms)允许元素迟到的时间超过 watermark 与窗口末端之间这个时长的元素将被丢弃。默认值为 0且仅对事件时间窗口有效L1866-L1875side_output_late_data(output_tag)把迟到数据送入指定OutputTag的侧输出1.16.0 引入。所谓迟到指watermark 已越过窗口末端 allowed_lateness。可用DataStream.get_side_output(tag)取回L1877-L1899迟到数据侧输出示例来自源码 docstring tag OutputTag(late-data, Types.TUPLE([Types.INT(), Types.STRING()])) main_stream ds.key_by(lambda x: x[1]) \ ... .window(TumblingEventTimeWindows.of(Time.seconds(5))) \ ... .side_output_late_data(tag) \ ... .reduce(lambda a, b: (a[0] b[0], b[1])) late_stream main_stream.get_side_output(tag)5.2 窗口计算函数reduce(reduce_function, window_functionNone, output_typeNone)先对窗口内元素做增量归约再可选交给WindowFunction或ProcessWindowFunction处理L1901-L1951。源码使用ReducingStateDescriptor(WINDOW_STATE_NAME, ...)WINDOW_STATE_NAME window-contents承载增量聚合状态。滚动时间窗口可做到每键只存一个元素滑动窗口按滑动粒度聚合每键每滑动间隔存一个元素自定义窗口可能无法增量聚合。 ds.key_by(lambda x: x[1]) \ ... .window(TumblingEventTimeWindows.of(Time.seconds(5))) \ ... .reduce(lambda a, b: (a[0] b[0], b[1]))aggregate(aggregate_function, window_functionNone, output_typeNone)基于AggregateFunction的增量聚合比 reduce 更灵活可维护累加器中间态。apply(window_function, output_typeNone)用WindowFunction对整窗口求值非增量。process(process_window_function, output_typeNone)用ProcessWindowFunction求值可访问窗口上下文与键。辅助方法get_execution_environment()、get_input_type()分别返回执行环境与窗口输入类型L1853-L1857。六、ConnectedStreams双流连接ConnectedStreams表示两个可能不同类型的流的连接data_stream.py。适合一个流上的操作直接影响另一个流的场景通常借助两流间共享的状态。源码文档给出的典型例子是动态规则流 数据流规则流提供规则数据流提供待处理元素连接算子把当前规则集合维护在状态中收到规则更新则更新状态收到数据元素则用当前规则处理它。概念上连接流可视为一个 Either 类型的并集流——它持有第一个流的类型或第二个流的类型。支持的方法L241-L244方法语义key_by(key_selector1, key_selector2)对两条连接流分别按键可用相同或不同选择器map(co_map_function, output_typeNone)用CoMapFunction联合转换分别处理两条流的元素flat_map(co_flat_map_function, output_typeNone)用CoFlatMapFunction联合转换process(co_process_function, output_typeNone)用CoProcessFunction联合处理可访问共享状态、定时器与侧输出rules env.from_collection([...], Types.ROW([...])) data env.from_collection([...], Types.ROW([...])) connected data.connect(rules) result connected.process(MyCoProcessFunction(), Types.ROW([...]))七、AllWindowedStream全局开窗流AllWindowedStream表示对整个非按键数据流按WindowAssigner切窗的流data_stream.py由DataStream.window_all()创建。与 WindowedStream 的差异在于元素按窗口分组而非按键所有元素进同一组窗口。若指定Evictor它会在 Trigger 触发之后、窗口实际计算之前剔除窗口内元素使用 evictor 会显著降低窗口性能因为无法使用窗口结果的预聚合。与 WindowedStream 类似它也是纯 API 构造运行时与窗口算子折叠为单个算子。方法集与 WindowedStream 完全一致trigger、allowed_lateness、side_output_late_data、reduce、aggregate、apply、process外加get_execution_environment、get_input_type。由于没有按键维度全局窗口的并行度天然受限通常为 1适用于需要全量数据参与的计算。八、BroadcastStream 与 BroadcastConnectedStream8.1 BroadcastStreamBroadcastStream由DataStream.broadcast(*MapStateDescriptor)创建data_stream.py。它把流中的元素广播到下游每个并行实例并隐式创建由MapStateDescriptor描述的BroadcastState。广播状态只允许从广播流一侧写入保证所有并行实例看到一致的状态。8.2 BroadcastConnectedStreamBroadcastConnectedStream表示键控或非键控流与BroadcastStream连接的结果data_stream.py。与 ConnectedStreams 类似它解决一端的操作直接影响另一端的问题但广播状态让规则在所有并行实例中可见因此可作用于另一条流的所有分区。典型的动态规则应用模式广播流携带规则并存入广播状态数据流携带待处理元素process方法内收到规则更新则写广播状态收到数据元素则从广播状态读取当前规则并应用。这是动态规则变更场景的官方推荐实现。支持的处理方法BroadcastConnectedStream.process(broadcast_process_function, output_typeNone)用BroadcastProcessFunction非键控连接时或KeyedBroadcastProcessFunction键控连接时处理L2642-L2657。broadcast_state_desc MapStateDescriptor(rules, Types.STRING(), Types.INT()) broadcast_stream rules_stream.broadcast(broadcast_state_desc) connected_stream data_stream.connect(broadcast_stream) result connected_stream.process(MyBroadcastProcessFunction(), Types.STRING())九、常见组合模式速查把上述类型串联起来就构成 PyFlink 流作业的完整骨架from pyflink.common import Types from pyflink.datastream import StreamExecutionEnvironment, OutputTag from pyflink.datastream.window import TumblingEventTimeWindows from pyflink.datastream.time import Time env StreamExecutionEnvironment.get_execution_environment() source env.from_collection([(1, a), (2, b), (1, c)], Types.TUPLE([Types.INT(), Types.STRING()])) # 1) 基础转换 键控 事件时间窗口 迟到数据侧输出 late_tag OutputTag(late, Types.TUPLE([Types.INT(), Types.INT()])) result source \ .map(lambda t: (t[0], 1), Types.TUPLE([Types.INT(), Types.INT()])) \ .key_by(lambda t: t[0]) \ .window(TumblingEventTimeWindows.of(Time.seconds(5))) \ .side_output_late_data(late_tag) \ .reduce(lambda a, b: (a[0], a[1] b[1])) late_data result.get_side_output(late_tag) result.name(window-sum).uid(window-sum-op).set_parallelism(2) result.print() # 2) execute_and_collect 拉回结果iterator 需关闭 with result.execute_and_collect(pyflink-demo) as it: for row in it: print(row)常见模式速查状态对齐与迁移为关键算子设置uid/set_uid_hash保证 savepoint 恢复与跨版本迁移。吞吐与延迟权衡用set_buffer_timeout(0)追求最低延迟用默认缓冲追求吞吐。资源隔离用slot_sharing_group把不同作业段隔离到不同槽。动态配置用 broadcast BroadcastConnectedStream.process实现不重启作业的规则热更新。结果复用用DataStream.cache()在有界流场景下跨作业复用中间结果。十、测试与验证参考仓库中的 test_data_stream.py 同时包含流模式DataStreamStreamingTests与批模式DataStreamBatchTests两组测试基类L784-L801覆盖 map/flat_map/filter/key_by/window/connect 等核心算子的端到端行为test_window.py 专门验证窗口分配与窗口函数test_slot_sharing_group.py 验证槽共享组的资源规格。阅读这些测试可以作为理解 API 语义与排查问题的第一手参考。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink DataStream API 全景参考从执行环境到连接器的完整编程模型指南PyFlink DataStream API 全景参考从执行环境到连接器的完整编程模型指南 本文以 Apache Flink 仓库中 PyFlink Data大数据流处理批处理数据工程Flink DataStream API 与 Table API 集成指南桥接转换、Changelog 流与统一批流执行Flink DataStream API 与 Table API 集成指南桥接转换、Changelog 流与统一批流执行 Flink 的 Table API大数据流处理批处理数据工程PyFlink Table Window 窗口 API 完全指南Tumble、Slide、Session 与 Over 窗口PyFlink Table Window 窗口 API 完全指南Tumble、Slide、Session 与 Over 窗口 窗口Window是流式数据处大数据流处理批处理数据工程上一篇nanowhale-100m革命性小型语言模型实现DeepSeek-V4架构的完整指南下一篇CANN/catlass TileCopyTlaL1 → L0A 偏特化创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网