☰
Apache Beam 多输出实战:用 ParDo 的 Side Output 将数据分流到多个 PCollection
2026/9/28 3:45:28 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache 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 Output:ParDo 的"一进多出"

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, main=None)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 output);
  • DoFnyieldpvalue.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, main=None, allow_unknown_tags=None):
  • *tags:非空时表示合法的 tag 白名单。若提供了白名单,之后在管道中访问未声明的 tag 会直接报错,从而尽早暴露拼写错误;
  • main=None:通过关键字参数main=...指定哪个 tag 作为主输出(不写时主输出默认是匿名输出);
  • allow_unknown_tags:允许访问未在*tags中声明的 tag(默认在声明了白名单时禁止)。

返回的对象是DoOutputsTuple(见 pvalue.py#L234),它支持三种访问方式:

  1. 下标访问results[tag];
  2. 属性访问results.tag(通过__getattr__实现,见 pvalue.py#L276-L281);
  3. 迭代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 给出了标准答案,下面逐段拆解。

定义两个输出 tag

num_below_100_tag = 'num_below_100' num_above_100_tag = 'num_above_100'

tag 是任意字符串,习惯上用描述性命名。这里一个表示"小于等于 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, main=num_below_100_tag)) results[num_below_100_tag] | 'Log numbers <= 100' >> beam.LogElements(prefix='Number <= 100: ') results[num_above_100_tag] | 'Log numbers > 100' >> beam.LogElements(prefix='Number > 100: ')

关键点:

  1. beam.Create([10, 50, 120, 20, 200, 0])构造输入 PCollection;
  2. beam.ParDo(ProcessNumbersDoFn()).with_outputs(num_above_100_tag, main=num_below_100_tag):
    • 声明附加输出 tagnum_above_100_tag;
    • 同时用main=num_below_100_tag把主输出命名为num_below_100_tag;
    • 返回results这个多输出元组;
  3. 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 的完整调用链,理解它有助于排查问题:

  1. ParDo.with_outputs(...)只是声明:在 core.py#L2373-L2377 中,它把附加 tag 存入self._extra_tags,把主 tag 存入self._main_tag,然后返回self,真正的 PCollection 尚未创建;
  2. 返回的DoOutputsTuple对象(pvalue.py#L234)在__init__中记录 pipeline、transform、tags 和 main_tag;
  3. 当你首次访问某个 tag(如results[tag])时,__getitem__(pvalue.py#L283-L324)才会真正创建对应的PCollection:
    • 对附加输出:调用self._transform.output_tags.add(tag)把 tag 注册为 ParDo 的真实输出,并创建带tag=tag的PCollection,同时把它挂到_MultiParDo及其内部 ParDo 的输出列表中(pvalue.py#L302-L316);
    • 对主输出:直接复用内部 ParDo 的匿名输出 PCollection(pvalue.py#L317-L322)。

也就是说,Beam 采用懒加载策略:只有下游真正消费某个输出分支时,该分支的 PCollection 才会被实例化并注册进执行图。这也解释了为什么"声明了 tag 但从未访问"不会产生副作用。

扩展应用与注意事项

多个附加输出

.with_outputs支持一次声明任意多个 tag,例如同时输出"错误日志"和"重试队列"两个分支:

results = (pcoll | beam.ParDo(MyDoFn()).with_outputs( 'errors', 'retry', main='main')) 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 抛TypeError(pvalue.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 目录下的课程)练习更多变换组合。

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

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:如何免费用浏览器读 EPUB:Epub.js Reader 实用上手指南
下一篇:如何使用 BallonsTranslator 一键翻译漫画

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询