- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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
相关推荐
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)是流式数据处
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考