大数据流处理批处理数据工程【免费下载链接】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.OutputTagPyFlink 通过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特别值得注意两个细节空 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。类型信息与序列化方式强相关不指定类型时走 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.pyfrom 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_infoTypes.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_typeTypes.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.pydef 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_typeoutput_tag.type_info)其要点如下在执行了process或窗口聚合等之后得到的DataStream上调用传入与发射侧输出时完全相同的OutputTag实例内部先通过getSideOutput拿到 JVM 侧的侧输出流再用一个map(lambda i: i, output_typeoutput_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_infoTypes.ROW([Types.STRING(), Types.INT()])) ds1 env.from_collection([(a, 0), (a, 8), (a, 4), (a, 6)], type_infoTypes.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.pyCoProcessFunctiontest_co_process_side_output两条输入流的交叉字段见 test_data_stream.pyBroadcastProcessFunctiontest_co_broadcast_side_output普通流与广播流元素的混合发射见 test_data_stream.pyKeyedProcessFunctiontest_keyed_process_side_output基于 keyed state 的累计值见 test_data_stream.pyKeyedCoProcessFunctiontest_keyed_co_process_side_output按 key 汇总的计数见 test_data_stream.py事件时间窗口test_side_output_late_data迟到数据见 test_window.py以多侧输出为例test_data_stream.pytag1 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_typeTypes.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_operatortest_data_stream.py验证了侧输出可以穿过上游算子链在ds.map(...).process(MyProcessFunction())之后调用get_side_output(tag)依然能取回数据说明侧输出与算子链式调用的兼容性良好。实现原理与序列化注意事项从源码层面看PyFlink 的侧输出链路包含三层协作Python 用户层用户构造OutputTag在过程函数中yield (tag, value)执行引擎层Python 侧的 UDF 执行器operations.py识别(OutputTag, value)元组调用SideOutputContext.collect按 tag id 路由数据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 状态一起被持久化。常见错误与最佳实践结合官方文档示例与源码校验逻辑使用侧输出时最容易踩的坑如下空 tag_idOutputTag()会立即抛出ValueErrorPython API 额外要求务必为每个侧输出起一个有语义的名字type_info 传错类型传入TypeInformation之外的对象如字符串会抛出TypeError请使用Types类中定义的工厂方法取流时 tag 不一致get_side_output必须使用发射侧输出时的同一个OutputTag实例或至少 tag_id 与 type_info 完全一致否则侧输出流为空且无任何报错忘记声明输出类型省略type_info时元素按 pickle 字节数组传输虽便捷但类型不透明跨语言互操作前建议显式声明版本限制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点击查看免费下载相关推荐Flink Side Outputs侧输出流完全指南定义 OutputTag、多类型分流与迟到数据处理Flink Side Outputs侧输出流完全指南定义 OutputTag、多类型分流与迟到数据处理 导读 在 Flink DataStream 开发中大数据流处理批处理数据工程Flink DataStream 旁路输出Side Output完全指南OutputTag 定义、Context 发射与 getSideOutput 消费Flink DataStream 旁路输出Side Output完全指南OutputTag 定义、Context 发射与 getSideOutput 消费大数据流处理批处理数据工程如何用go2rtc构建零延迟的智能摄像头流媒体系统5个实战配置技巧如何用go2rtc构建零延迟的智能摄像头流媒体系统5个实战配置技巧 go2rtc是一款跨平台的摄像头流媒体解决方案专为中级用户和开发者设计。它支持RTSP、大数据流处理批处理数据工程上一篇jc net_localgroup 解析器将 Windows net localgroup 输出转为结构化 JSON 的完整指南下一篇Warden Protocol 治理机制去中心化决策的完整实现创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
wp-calypso 结账表单校验与输入掩码:getCreditCardType 与 maskField 实战解析 前端CMS 【免费下载链接】wp-calypso The JavaScript and API powered WordPress.com 项目地址: https://gitcode.com/gh_mirrors/wp/wp-calypso 点击查看 免费下载 本指南以 client/lib/checkout/README.md 为核心,深入讲解 WordPress.com 前端项目 wp… · 2026/9/25 3:02:00
ClawHub 本地开发指南:从环境搭建、Worktree 快路径到 PR 提交门禁的完整实践 后端前端AI 技能AI 插件搜索引擎 【免费下载链接】clawhub Skill Plugin Registry for OpenClaw 项目地址: https://gitcode.com/gh_mirrors/mo/clawhub 点击查看 免费下载 本文以 ClawHub(OpenClaw 的公开技能与插件注册表)仓库的 CONTRIB… · 2026/9/25 3:01:54
IDURAR ERP/CRM 功能全景解析:基于 MERN 技术栈的开源企业资源规划与客户关系管理软件 后端前端企业应用CRM 【免费下载链接】idurar-erp-crm Free Open Source ERP CRM Software Accounting Invoicing | Node.Js React 项目地址: https://gitcode.com/gh_mirrors/id/idurar-erp-crm 点击查看 免费下载 IDURAR 是一款免费开源的 ERP(企业资… · 2026/9/25 3:01:54
BLE Mesh抓包实战:nRF52840 Dongle与Wireshark联合调试指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:45
Android Studio无数据库注册页实战:从布局到传参的完整指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:45
沪深主板打板胜率统计:数据清洗与实战规则 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:45
Zabbix与Prometheus监控选型实战:从场景DNA到工具链协同 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:45
甘肃职称评审论文查重系统全解析:从检测逻辑到避坑指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:44
Spring Boot校园闲置交易毕设:表结构、订单状态机与避坑指南 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 4:59:38
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37