news 2026/9/25 2:25:09

PyFlink DataStream API 完全指南:从基础流转换到窗口、连接与广播流

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
PyFlink DataStream API 完全指南:从基础流转换到窗口、连接与广播流
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】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_group

  • set_parallelism(parallelism):设置当前算子的并行度(L139-L147)。
  • set_max_parallelism(max_parallelism):设置算子最大并行度,它定义动态扩缩容的上限,也决定分区状态(partitioned state)的 key group 数量(L149-L160)。
  • force_non_parallel():把并行度与最大并行度都强制为 1,且后续不可再设为非 1(L183-L191)。适用于connect等要求单实例的操作。
  • slot_sharing_group(slot_sharing_group):设置槽共享组。同一组内的并行实例会被尽可能放置到同一个 TaskManager 槽中;未显式设置时继承输入算子的组;可通过传入字符串'default'或 SlotSharingGroup 对象指定(源码 L234-L253 支持字符串与带资源规格的对象两种形式)。

1.3 运行时行为调优:set_buffer_timeout / start_new_chain / disable_chaining

  • set_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、ProcessFunction):

  • map(func, output_type=None):每个元素恰好输出一个元素(L273-L314)。源码内部把它包装成MapProcessFunctionAdapter,复用process通道实现。注意:不指定output_type时,输出数据会以 pickle 原始字节数组的形式序列化。
  • flat_map(func, output_type=None):每个元素可输出任意数量(含 0 个)元素(L315-L355)。
  • filter(func):仅保留函数返回True的元素(L415-L455)。
  • process(func, output_type=None):最底层的转换原语,每个元素可输出 0 个或多个元素,并能访问上下文(L641-L664)。上述 map/flat_map/filter 均是它的语法糖封装。

1.5 键控与分区:key_by / partition_custom / 分区器家族

  • key_by(key_selector, key_type=None):按指定键把流分区为 KeyedStream(L357-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-L534)
rescale()轮询分发到下游的子集实例,子集大小取决于上下游并行度的倍数关系(L555-L573)
rebalance()轮询分发到下游全部实例(L575-L582)
forward()元素转发到下游的本地子任务(L584-L591)
broadcast()元素广播到下游每个并行实例(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_watermarks

  • window_all(window_assigner):对非按键分组的流开窗,返回 AllWindowedStream(1.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 / connect

  • union(*streams):将多个同类型DataStream 合并(L473-L493)。源码会自动处理KeyedStream参数(取_values()后再合并)。
  • connect(ds):连接两个(可能不同类型的)流,返回 ConnectedStreams;若ds是 BroadcastStream 则返回 BroadcastConnectedStream(1.16.0 起支持,L495-L525)。

1.8 输出与执行:add_sink / sink_to / print / execute_and_collect

  • add_sink(sink_func):为流添加SinkFunction型 sink(L811-L819)。
  • sink_to(sink):添加新版Sink接口的 sink(如 Kafka/Pulsar 等 connectors 模块中的实现),并支持SupportsPreprocessing预处理转换(L821-L837)。
  • print(sink_identifier=None):把流写到标准输出(L869-L884)。注意:它输出到运行代码的机器(即 Flink worker)的 stdout,且不具备容错能力。内部会先调用_align_output_type把 pickle 对象或 Row 转成可读字符串。
  • execute_and_collect(job_execution_name=None, limit=None):触发分布式执行并通过 Flink REST API 把事件拉回当前进程(L839-L867)。不传limit时返回CloseableIterator,必须关闭以释放集群资源;传limit时返回 Python list。执行前会自动调用PythonOperatorChainingOptimizer做算子链接优化(L913-L924)。

1.9 侧输出与缓存:get_side_output / cache

  • get_side_output(output_tag):取出算子通过指定OutputTag发出的旁路数据流(1.16.0 引入,L886-L897)。通常与WindowedStream.side_output_late_data配合收集迟到数据。
  • cache():缓存变换的中间结果,返回 CachedDataStream(1.16.0 引入,L899-L911)。仅支持有界流,当前仅支持 block 模式;缓存在该中间结果首次被计算时惰性生成,并在StreamExecutionEnvironment关闭时清除。

1.10 类型与执行环境

  • get_type():返回流的TypeInformation(L162-L168),所有output_type参数都可借此推导。
  • get_execution_environment()/get_execution_config():返回创建该流的执行环境及其配置(L170-L181)。

二、KeyedStream:按键分区流

KeyedStream表示算子状态按键分区(通过KeySelector)的 DataStream(data_stream.py)。除了shuffle、forward、key_by等分区类操作不可用外,普通 DataStream 的转换都可使用,且reduce、sum等"归约类"操作只作用于相同键的元素。

2.1 键控转换与归约

  • map/flat_map/filter/process:与普通流等价,但内部使用KeyedProcessFunction适配器包装,可访问键控上下文(L1113-L1195、L1616)。
  • reduce(func):对每个键分组内元素做归约(L1197)。
  • sum(position=0)、min(position=0)、max(position=0)、min_by(position=0)、max_by(position=0):按字段位置(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):按键开窗,返回 WindowedStream(L1644)。
  • count_window(size, slide=0):按元素计数开窗(L1658),slide=0时退化为滚动计数窗口。窗口分配器定义在 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)直接指定 JobVertexID
set_parallelism(parallelism)设置 sink 并行度
set_description(description)设置详细描述(1.15.0 引入,用于 JSON plan 与 Web UI)
disable_chaining()关闭该算子的链接优化
slot_sharing_group(group)设置槽共享组(字符串或SlotSharingGroup对象)

四、CachedDataStream:可复用缓存流

CachedDataStream表示中间结果会被缓存的 DataStream(data_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-L1843):WindowedStream 纯属 API 层面的构造,运行时它会被与 KeyedStream 和窗口算子折叠为单个算子执行——这正是窗口聚合性能的根源。

5.1 窗口行为配置

方法语义
trigger(trigger)设置触发窗口发射的Trigger(L1859-L1864)。不设置时使用 WindowAssigner 的默认 Trigger
allowed_lateness(time_ms)允许元素迟到的时间;超过 watermark 与窗口末端之间这个时长的元素将被丢弃。默认值为 0,且仅对事件时间窗口有效(L1866-L1875)
side_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_function=None, output_type=None):先对窗口内元素做增量归约,再(可选)交给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_function=None, output_type=None):基于AggregateFunction的增量聚合,比 reduce 更灵活(可维护累加器中间态)。

  • apply(window_function, output_type=None):用WindowFunction对整窗口求值(非增量)。

  • process(process_window_function, output_type=None):用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_type=None)用CoMapFunction联合转换(分别处理两条流的元素)
flat_map(co_flat_map_function, output_type=None)用CoFlatMapFunction联合转换
process(co_process_function, output_type=None)用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 与 BroadcastConnectedStream

8.1 BroadcastStream

BroadcastStream由DataStream.broadcast(*MapStateDescriptor)创建(data_stream.py)。它把流中的元素广播到下游每个并行实例,并隐式创建由MapStateDescriptor描述的BroadcastState。广播状态只允许从广播流一侧写入,保证所有并行实例看到一致的状态。

8.2 BroadcastConnectedStream

BroadcastConnectedStream表示"键控或非键控流"与BroadcastStream连接的结果(data_stream.py)。与 ConnectedStreams 类似,它解决"一端的操作直接影响另一端"的问题,但广播状态让规则在所有并行实例中可见,因此可作用于另一条流的所有分区。

典型的动态规则应用模式:广播流携带规则并存入广播状态;数据流携带待处理元素;process方法内,收到规则更新则写广播状态,收到数据元素则从广播状态读取当前规则并应用。这是"动态规则变更"场景的官方推荐实现。

支持的处理方法:

  • BroadcastConnectedStream.process(broadcast_process_function, output_type=None):用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
点击查看免费下载
上一篇:nanowhale-100m:革命性小型语言模型实现DeepSeek-V4架构的完整指南
下一篇:CANN/catlass TileCopyTla(L1 → L0A 偏特化)

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 2:24:59

applera1n工作原理揭秘:checkm8硬件漏洞如何支撑整个iCloud Bypass

applera1n工作原理揭秘:checkm8硬件漏洞如何支撑整个iCloud Bypass 【免费下载链接】applera1n icloud bypass for ios 15-16 项目地址: https://gitcode.com/gh_mirrors/ap/applera1n 还在为iPhone被iCloud账号锁住而烦恼?本文带你从底层原理讲透…

作者头像 李华
网站建设 2026/9/25 2:24:27

LinuxKit 内核性能分析指南:使用 perf 工具进行容器系统性能剖析

操作系统云原生容器运行时 【免费下载链接】linuxkit A toolkit for building secure, portable and lean operating systems for containers 项目地址: https://gitcode.com/gh_mirrors/li/linuxkit 点击查看 免费下载 导读 本文介绍如何在 LinuxKit 构建的最小化…

作者头像 李华
网站建设 2026/9/25 2:22:34

Windows原生命令certutil计算MD5的原理与实战

1. 为什么在Windows下必须亲手算MD5?不是有图形工具吗?你有没有遇到过这种情况:下载完一个ISO镜像,官网只给了MD5校验值,你双击打开某个“MD5计算器.exe”,拖进去文件,结果弹窗提示“无法读取文…

作者头像 李华