news 2026/9/25 3:02:04

PyFlink DataStream Side Outputs 完全指南:使用 OutputTag 处理侧输出流

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
PyFlink DataStream Side Outputs 完全指南:使用 OutputTag 处理侧输出流
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

导读

在 PyFlink DataStream API 中,侧输出(Side Outputs)是一种从主数据流中"分流"出额外结果流的能力。本文聚焦于pyflink.datastream包中的核心类型OutputTag(见 sideoutput.rst),系统讲解其三种构造方式、类型推断规则,以及如何配合ProcessFunction的yield机制、DataStream.get_side_output和窗口的side_output_late_data实现迟到数据处理、异常事件隔离等典型场景。读完本文,你将能够熟练地为一个算子声明多个带类型的侧输出通道,并在主、侧两条数据流上分别做后续处理。

OutputTag:侧输出的"身份证"

在 Flink 的流处理模型中,一个算子除了产生一条主输出流(main output)之外,还可以通过侧输出(side output)向任意数量的附加数据流发射数据。侧输出与主输出互不影响,非常适合处理以下场景:

  • 迟到数据:窗口已经触发计算后,迟到的数据不再进入主结果流,而是被引导到专门的侧输出流,供下游做延迟补偿或监控;
  • 异常/告警事件:在正常业务数据之外,将校验失败、格式错误或超阈值的事件单独输出,交给独立的告警或修复链路;
  • 多路分流:一条输入流经过一个算子处理后,按不同规则拆分成多条不同类型的数据流。

而OutputTag正是这条侧输出通道的"身份证"——它用**名字(tag_id)标识通道,用类型信息(type_info)**约束通道中元素的类型。其定义位于 output_tag.py,类签名如下:

class OutputTag(object): def __init__(self, tag_id: str, type_info: Optional[Union[TypeInformation, list]] = None): ...

对应的 Java 底层类型是org.apache.flink.util.OutputTag,PyFlink 通过get_java_output_tag()方法借助 Java Gateway 将其桥接为 JVM 对象(见 output_tag.py)。

三种构造方式与类型推断规则

根据 sideoutput.rst 的官方示例,OutputTag支持三种构造方式,分别对应不同的类型信息来源:

# 方式一:显式指定输出类型 >>> info = OutputTag("late-data", Types.TUPLE([Types.STRING(), Types.LONG()])) # 方式二:隐式将 list 包装为 Types.ROW >>> info_row = OutputTag("row", [Types.STRING(), Types.LONG()]) # 方式三:隐式使用 pickle 序列化 >>> info_side = OutputTag("side") # 错误示例:tag id 不能为空字符串(Python API 的额外要求) >>> info_error = OutputTag("")

三种方式的底层判定逻辑位于 output_tag.py,可以归纳为一张决策表:

传入的type_info参数实际生效的类型信息说明
None(省略)Types.PICKLED_BYTE_ARRAY()元素以 Python pickle 序列化后的字节数组传输,最灵活但开销最大
list,如[Types.STRING(), Types.LONG()]RowTypeInfo(type_info)列表中的每个TypeInformation自动被包装为一行 ROW 的字段类型
TypeInformation子类实例原样使用例如Types.INT()、Types.TUPLE(...)等,类型最精确
其他类型抛出TypeError构造函数会校验:"OutputTag type_info must be None, list or TypeInformation"

特别值得注意两个细节:

  1. 空 tag_id 是硬性错误:__init__中会执行if not tag_id: raise ValueError("OutputTag tag_id cannot be None or empty string")。这是 Python API 相对于 Java API 的额外约束——Java 侧允许空字符串 tag,但 PyFlink 出于可读性和调试友好性禁止了它(见 output_tag.py)。
  2. 类型信息与序列化方式强相关:不指定类型时走 pickle 序列化,这意味着侧输出流中元素的 Java 侧类型是byte[],后续如果想在 Java/SQL 侧复用该流,需要显式声明类型以获得确定性更高的编码。

在 ProcessFunction 中发射侧输出:yield (tag, value)

PyFlink 的过程函数(ProcessFunction、KeyedProcessFunction、CoProcessFunction、BroadcastProcessFunction等)使用 Python 生成器(generator)语义:在process_element等方法中,直接yield元素会进入主输出流,yield (output_tag, value)二元组则会把value发射到该 tag 对应的侧输出流。

底层实现见 operations.py:执行引擎对函数产出的每个值做检查,当isinstance(value, tuple) and isinstance(value[0], OutputTag)时,就调用side_output_context.collect(output_tag.tag_id, value[1])将其路由到侧输出。

一个完整的ProcessFunction侧输出示例(对应官方测试 test_data_stream.py):

from pyflink.common.typeinfo import Types from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.functions import ProcessFunction from pyflink.datastream.output_tag import OutputTag env = StreamExecutionEnvironment.get_execution_environment() tag = OutputTag("side", Types.INT()) ds = env.from_collection([('a', 0), ('b', 1), ('c', 2)], type_info=Types.ROW([Types.STRING(), Types.INT()])) class MyProcessFunction(ProcessFunction): def process_element(self, value, ctx: 'ProcessFunction.Context'): yield value[0] # 进入主输出流 yield tag, value[1] # 进入 "side" 侧输出流 ds2 = ds.process(MyProcessFunction(), output_type=Types.STRING()) main_sink = ... # 自定义 Sink ds2.add_sink(main_sink) side_sink = ... ds2.get_side_output(tag).add_sink(side_sink) env.execute("test_process_side_output")

该测试的断言结果清晰地体现了主、侧分流效果:

  • 主输出流收到['a', 'b', 'c'](即每行的第一个字段);
  • 侧输出流收到['0', '1', '2'](即每行的第二个整数字段)。

读取侧输出流:DataStream.get_side_output(tag)

发射进侧输出的元素不会自动出现在下游,必须通过get_side_output显式取回,该方法定义在 data_stream.py:

def get_side_output(self, output_tag: OutputTag) -> 'DataStream': ds = DataStream(self._j_data_stream.getSideOutput(output_tag.get_java_output_tag())) return ds.map(lambda i: i, output_type=output_tag.type_info)

其要点如下:

  • 在执行了process(或窗口聚合等)之后得到的DataStream上调用,传入与发射侧输出时完全相同的OutputTag实例;
  • 内部先通过getSideOutput拿到 JVM 侧的侧输出流,再用一个map(lambda i: i, output_type=output_tag.type_info)将元素转换回 Python 对象——因此侧输出流的元素类型由OutputTag声明时的type_info决定;
  • 该方法自1.16.0版本加入(docstring 中标有.. versionadded:: 1.16.0),使用旧版本时请升级 PyFlink 或改用其他取流方式;
  • 取流与发射必须使用同一个 tag,否则拿到的将是空流,且不会报错。

窗口场景:用side_output_late_data收留迟到数据

窗口迟到数据是侧输出最经典的应用。WindowedStream与AllWindowedStream都提供了side_output_late_data(output_tag)方法(见 data_stream.py),用于把在"窗口结束时间 + allowed_lateness"之后才到达的数据引导到指定侧输出:

>>> 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)

这段官方示例中值得注意的编排顺序:.side_output_late_data(tag)必须在窗口算子(reduce/process/aggregate等)之前调用,然后对窗口算子返回的DataStream调用get_side_output(tag)即可取回迟到数据流。

窗口测试用例 test_window.py 提供了完整可运行验证:

tag = OutputTag('late-data', type_info=Types.ROW([Types.STRING(), Types.INT()])) ds1 = env.from_collection([('a', 0), ('a', 8), ('a', 4), ('a', 6)], type_info=Types.ROW([Types.STRING(), Types.INT()])) ds2 = ds1.assign_timestamps_and_watermarks(watermark_strategy) \ .key_by(lambda e: e[0]) \ .window(TumblingEventTimeWindows.of(Time.milliseconds(5))) \ .allowed_lateness(0) \ .side_output_late_data(tag) \ .process(CountWindowProcessFunction(), Types.TUPLE([Types.STRING(), Types.LONG(), Types.LONG(), Types.INT()])) ds2.add_sink(main_sink) ds2.get_side_output(tag).add_sink(side_sink) env.execute('test_side_output_late_data')

测试最终断言:

  • 主输出流只有完整落入窗口区间的数据:['(a,0,5,1)', '(a,5,10,2)'];
  • 侧输出流捕获了迟到的('a', 4):['+I[a, 4]']。

这一用例同时印证了 sideoutput.rst 中OutputTag("late-data", ...)命名的惯例:用 tag 的名字表达数据的语义(late-data),用类型信息约束数据的结构。

多 tag 与多算子类型的组合应用

侧输出并非ProcessFunction的专利。仓库测试覆盖了多种算子场景,均可通过yield (tag, value)与get_side_output组合使用:

算子类型测试用例侧输出内容
ProcessFunctiontest_process_side_output/test_process_multiple_side_output每条元素的整型字段;可同时声明tag1、tag2两条侧输出(见 test_data_stream.py)
CoProcessFunctiontest_co_process_side_output两条输入流的交叉字段(见 test_data_stream.py)
BroadcastProcessFunctiontest_co_broadcast_side_output普通流与广播流元素的混合发射(见 test_data_stream.py)
KeyedProcessFunctiontest_keyed_process_side_output基于 keyed state 的累计值(见 test_data_stream.py)
KeyedCoProcessFunctiontest_keyed_co_process_side_output按 key 汇总的计数(见 test_data_stream.py)
事件时间窗口test_side_output_late_data迟到数据(见 test_window.py)

以多侧输出为例(test_data_stream.py):

tag1 = OutputTag("side1", Types.INT()) tag2 = OutputTag("side2", Types.STRING()) class MyProcessFunction(ProcessFunction): def process_element(self, value, ctx: 'ProcessFunction.Context'): yield value[0] yield tag1, value[1] yield tag2, value[0] + str(value[1]) ds2 = ds.process(MyProcessFunction(), output_type=Types.STRING()) ds2.get_side_output(tag1).add_sink(side1_sink) ds2.get_side_output(tag2).add_sink(side2_sink)

输入[('a', 0), ('b', 1), ('c', 2)]时,side1收到['0', '1', '2'],side2收到['a0', 'b1', 'c2']——可见每个 tag 拥有独立的元素类型与数据通道,互不干扰。

此外,test_side_output_chained_with_upstream_operator(test_data_stream.py)验证了侧输出可以穿过上游算子链:在ds.map(...).process(MyProcessFunction())之后调用get_side_output(tag)依然能取回数据,说明侧输出与算子链式调用的兼容性良好。

实现原理与序列化注意事项

从源码层面看,PyFlink 的侧输出链路包含三层协作:

  1. Python 用户层:用户构造OutputTag,在过程函数中yield (tag, value);
  2. 执行引擎层:Python 侧的 UDF 执行器(operations.py)识别(OutputTag, value)元组,调用SideOutputContext.collect按 tag id 路由数据;
  3. JVM 层:get_java_output_tag()通过 Java Gateway 创建org.apache.flink.util.OutputTag(tag_id, type_info.get_java_type_info()),最终与 Flink 原生的getSideOutputAPI 打通(见 output_tag.py)。

需要特别提醒的序列化陷阱:OutputTag内部缓存的 Java 对象(_j_output_tag)不能被 Python pickle 直接序列化,因此在算子状态快照或跨进程传递场景下,output_tag.py 通过自定义__getstate__/__setstate__主动清空 Java 引用,只保留(tag_id, type_info)两个纯 Python 字段,反序列化后惰性重建 Java 对象。这保证了OutputTag可以安全地随 UDF 状态一起被持久化。

常见错误与最佳实践

结合官方文档示例与源码校验逻辑,使用侧输出时最容易踩的坑如下:

  1. 空 tag_id:OutputTag("")会立即抛出ValueError(Python API 额外要求),务必为每个侧输出起一个有语义的名字;
  2. type_info 传错类型:传入TypeInformation之外的对象(如字符串)会抛出TypeError,请使用Types类中定义的工厂方法;
  3. 取流时 tag 不一致:get_side_output必须使用发射侧输出时的同一个OutputTag实例(或至少 tag_id 与 type_info 完全一致),否则侧输出流为空且无任何报错;
  4. 忘记声明输出类型:省略type_info时元素按 pickle 字节数组传输,虽便捷但类型不透明,跨语言互操作前建议显式声明;
  5. 版本限制:DataStream.get_side_output从 1.16.0 起可用,side_output_late_data需要配合allowed_lateness(默认 0)才能捕获迟到数据。

推荐的工程实践是:在作业顶层集中定义所有OutputTag常量,用语义化命名(如"late-data"、"invalid-events")配合精确的类型信息,既便于get_side_output复用,也让侧输出流在 Flink Web UI 与监控指标中一目了然。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:jc net_localgroup 解析器:将 Windows `net localgroup` 输出转为结构化 JSON 的完整指南
下一篇:Warden Protocol 治理机制:去中心化决策的完整实现

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

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

F´ 飞行软件框架安装指南:环境准备、工具链部署与故障排查

嵌入式系统编程 【免费下载链接】fprime F - A flight software and embedded systems framework 项目地址: https://gitcode.com/gh_mirrors/fpri/fprime 点击查看 免费下载 本指南面向想要在 Linux 或 macOS 上快速搭建 F(F Prime)飞行软件…

作者头像 李华
网站建设 2026/9/25 3:01:08

OpenShift Origin 容器化部署与 Sample App 环境准备指南

测试云原生质量保障 【免费下载链接】origin Conformance test suite for OpenShift 项目地址: https://gitcode.com/gh_mirrors/or/origin 点击查看 免费下载 本文基于 origin 仓库中的 container-setup.md 展开,介绍如何以 Docker 容器方式拉起一个自…

作者头像 李华
网站建设 2026/9/25 3:00:18

PySide6+PyInstaller实战:搞怪小程序桌面开发与打包

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华