ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Beam 多输出实战:用 ParDo 的 Side Output 将数据分流到多个 PCollection

Apache Beam 多输出实战:用 ParDo 的 Side Output 将数据分流到多个 PCollection 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的ParDo是构建数据处理管线的核心变换但大多数教程只介绍它返回单个主输出main output。当业务需要一个元素同时归属多个分支例如把大于 100 的数与小于等于 100 的数拆开处理时就需要用到 Side Output附加输出机制。本文基于 Apache Beam Python SDK 的官方 Kata 练习 learning/katas/python/Core Transforms/Side Output/Side Output/task.md完整讲解pvalue.TaggedOutput与.with_outputs的用法、底层实现与测试验证方式帮助你掌握在单个DoFn内输出多个 PCollection 的标准写法。什么是 Side OutputParDo 的一进多出ParDo变换对每个输入元素执行用户自定义的DoFn处理逻辑通常只产生一个主输出 PCollection——也就是apply或管道操作返回的那个集合。然而在很多真实场景中单个DoFn需要同时产出多类结果按数值阈值分流如本 Kata 的大于/小于等于 100把正常数据与异常数据分开输出异常走单独的告警分支一条数据同时触发多条下游分支例如既做汇总又做明细落库。Beam 的设计是ParDo始终有一个主输出同时可以携带任意数量的附加输出Side Output。当声明了多个输出时ParDo会返回一个把所有输出包括主输出捆绑在一起的对象后续代码再按 tag标签逐一取出对应的 PCollection。Kata 的任务原文明确写道While ParDo always produces a main output PCollection (as the return value from apply), you can also have your ParDo produce any number of additional output PCollections. If you choose to have multiple outputs, your ParDo returns all of the output PCollections (including the main output) bundled together.Kata 目标为你的ParDo实现附加输出把大于 100 的数字单独输出到一个分支。核心 API 速查TaggedOutput 与 with_outputs实现 Side Output 只需要两个配套 API它们都在任务提示中明确点名API位置作用pvalue.TaggedOutput(tag, value)sdks/python/apache_beam/pvalue.py在DoFn.process内包装元素指定它发往哪个带 tag 的输出.with_outputs(*tags, mainNone)sdks/python/apache_beam/transforms/core.py声明 ParDo 的附加输出 tag 列表返回可按下标/属性访问的多输出元组TaggedOutput 的语义在源码 sdks/python/apache_beam/pvalue.py#L331-L341 中TaggedOutput的文档注释解释了它的设计意图ParDo, Map, and FlatMap transforms can emit values on multiple outputs which are distinguished by string tags. The DoFn will return plain values if it wants to emit on the main output and TaggedOutput objects if it wants to emit a value on a specific tagged output.也就是说规则非常简单DoFn直接 yield 普通值→ 元素进入主输出main outputDoFnyieldpvalue.TaggedOutput(tag, value)→ 元素进入名为tag的附加输出。TaggedOutput的构造函数还会对非字符串 tag 抛出TypeError见 pvalue.py#L339-L341保证 tag 一定是字符串。with_outputs 的参数在 sdks/python/apache_beam/transforms/core.py#L1806-L1830 中with_outputs的签名与行为如下def with_outputs(self, *tags, mainNone, allow_unknown_tagsNone):*tags非空时表示合法的 tag 白名单。若提供了白名单之后在管道中访问未声明的 tag 会直接报错从而尽早暴露拼写错误mainNone通过关键字参数main...指定哪个 tag 作为主输出不写时主输出默认是匿名输出allow_unknown_tags允许访问未在*tags中声明的 tag默认在声明了白名单时禁止。返回的对象是DoOutputsTuple见 pvalue.py#L234它支持三种访问方式下标访问results[tag]属性访问results.tag通过__getattr__实现见 pvalue.py#L276-L281迭代for pcoll in results:先主输出后附加输出见 pvalue.py#L269-L274。当通过__getitem__访问某个 tag 时若该 tag 既不是主输出也不在白名单中且未开启allow_unknown_tags会抛出ValueError见 pvalue.py#L292-L295这是 Beam 帮你防呆的机制。完整可运行示例按 100 阈值分流任务目录 learning/katas/python/Core Transforms/Side Output/Side Output/task.py 给出了标准答案下面逐段拆解。定义两个输出 tagnum_below_100_tag num_below_100 num_above_100_tag num_above_100tag 是任意字符串习惯上用描述性命名。这里一个表示小于等于 100 的分支一个表示大于 100 的分支。在 DoFn 内用 TaggedOutput 分流class ProcessNumbersDoFn(beam.DoFn): def process(self, element): if element 100: yield element # 普通值 → 主输出 else: yield pvalue.TaggedOutput(num_above_100_tag, element) # → 附加输出注意这里巧妙的组合主输出并不匿名而是通过with_outputs(main...)显式命名为num_below_100。也就是说小于等于 100的数据走主输出直接yield普通值即可大于 100的数据走附加输出必须用TaggedOutput包装。这正是主输出 附加输出的典型分工。声明附加输出并同时消费两个分支with beam.Pipeline() as p: results \ (p | beam.Create([10, 50, 120, 20, 200, 0]) | beam.ParDo(ProcessNumbersDoFn()) .with_outputs(num_above_100_tag, mainnum_below_100_tag)) results[num_below_100_tag] | Log numbers 100 beam.LogElements(prefixNumber 100: ) results[num_above_100_tag] | Log numbers 100 beam.LogElements(prefixNumber 100: )关键点beam.Create([10, 50, 120, 20, 200, 0])构造输入 PCollectionbeam.ParDo(ProcessNumbersDoFn()).with_outputs(num_above_100_tag, mainnum_below_100_tag)声明附加输出 tagnum_above_100_tag同时用mainnum_below_100_tag把主输出命名为num_below_100_tag返回results这个多输出元组results[tag]分别取出两个分支各接一个beam.LogElements打印日志。预期运行输出为Number 100: 10 Number 100: 50 Number 100: 20 Number 100: 0 Number 100: 120 Number 100: 200元素顺序由 runner 决定不保证与输入顺序一致。测试用例如何验证多输出正确性Kata 附带的测试 learning/katas/python/Core Transforms/Side Output/Side Output/tests/test_task.py 清晰地展示了验收标准numbers_below_100 [0, 10, 20, 50] numbers_above_100 [120, 200] answers [] for num in numbers_below_100: answers.append(Number 100: num) for num in numbers_above_100: answers.append(Number 100: num) for ans in answers: self.assertIn(ans, output, Incorrect output. Output the numbers to the output tags accordingly.)测试从运行输出中抓取task.py的日志断言6 个输入元素中0、10、20、50四个元素必须出现在Number 100分支120、200两个元素必须出现在Number 100分支。这验证了附加输出只接收大于 100 的数据、主输出接收其余数据的分流正确性。在 task-info.yaml 中本练习被标记为complexity: BASIC分类为Filtering与Multiple Outputs即过滤 多输出两个能力点的组合训练。底层原理DoOutputsTuple 与 tag 的绑定过程从源码可以看到 Side Output 的完整调用链理解它有助于排查问题ParDo.with_outputs(...)只是声明在 core.py#L2373-L2377 中它把附加 tag 存入self._extra_tags把主 tag 存入self._main_tag然后返回self真正的 PCollection 尚未创建返回的DoOutputsTuple对象pvalue.py#L234在__init__中记录 pipeline、transform、tags 和 main_tag当你首次访问某个 tag如results[tag]时__getitem__pvalue.py#L283-L324才会真正创建对应的PCollection对附加输出调用self._transform.output_tags.add(tag)把 tag 注册为 ParDo 的真实输出并创建带tagtag的PCollection同时把它挂到_MultiParDo及其内部 ParDo 的输出列表中pvalue.py#L302-L316对主输出直接复用内部 ParDo 的匿名输出 PCollectionpvalue.py#L317-L322。也就是说Beam 采用懒加载策略只有下游真正消费某个输出分支时该分支的 PCollection 才会被实例化并注册进执行图。这也解释了为什么声明了 tag 但从未访问不会产生副作用。扩展应用与注意事项多个附加输出.with_outputs支持一次声明任意多个 tag例如同时输出错误日志和重试队列两个分支results (pcoll | beam.ParDo(MyDoFn()).with_outputs( errors, retry, mainmain)) errors results[errors] retry results[retry] main results[main]属性访问DoOutputsTuple支持results.num_above_100这样的属性式访问pvalue.py#L276-L281代码更简洁但注意 tag 需是合法 Python 标识符。注意事项主输出只有一个无论附加输出有多少个ParDo的主输出始终只有一个。若想让某个分支成为主输出用main关键字重命名即可tag 必须是字符串TaggedOutput构造时会对非字符串 tag 抛TypeErrorpvalue.py#L339-L341尽早声明白名单在with_outputs(*tags)中列出所有合法 tag可以借助ValueError在构建阶段拦截拼错的 tag而不是等运行时才发现不要把 Side Output 与 Side Input 混淆Side Input 是把额外 PCollection 作为输入广播进DoFn用AsSingleton/AsList等而 Side Output 是让DoFn同时产出多个输出 PCollection二者方向相反。小结Side Output 是 Apache Beam 多分支数据处理的基础能力本 Kata 用按 100 阈值分流的最小示例展示了它的完整用法DoFn内直接yield普通值进主输出yield pvalue.TaggedOutput(tag, value)进附加输出ParDo(...).with_outputs(附加tag, main主tag)声明并取回捆绑的多输出元组用results[tag]分别消费每个分支配套测试通过日志断言逐元素验证分流结果tests/test_task.py。掌握这一模式后无论是正常/异常数据分流、多阈值分桶还是同源多下游都能用同一套 API 干净地实现。想继续巩固可前往本仓库的其他 Python Kata如 learning/katas/python/Core Transforms 目录下的课程练习更多变换组合。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java 实战用 ParDo 多路输出Side Output拆分大于 100 的数字流Apache Beam Java 实战用 ParDo 多路输出Side Output拆分大于 100 的数字流 Apache Beam 的 ParDo 变大数据批处理流处理数据工程Apache Beam Kotlin Katas 实战用 ParDo 的 Side Output额外输出实现数据分流Apache Beam Kotlin Katas 实战用 ParDo 的 Side Output额外输出实现数据分流 本指南以 Apache Beam 仓大数据批处理流处理数据工程Apache Beam Kotlin 实战使用 ParDo 与 MultiOutputReceiver 实现多输出Side OutputKata 全解Apache Beam Kotlin 实战使用 ParDo 与 MultiOutputReceiver 实现多输出Side OutputKata 全解 A大数据批处理流处理数据工程上一篇如何免费用浏览器读 EPUBEpub.js Reader 实用上手指南下一篇如何使用 BallonsTranslator 一键翻译漫画创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表