- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
Side inputs(侧输入)是 Apache Beam 中ParDo变换在主输入之外接收的附加数据源,它让每个元素的处理逻辑可以在运行时动态读取外部数据,而无需预先写死常量。本指南以仓库文档 29_advanced_side_inputs.md 为骨架,结合 Python 与 Java SDK 的真实源码与示例,系统讲解 side inputs 的声明方式、五种视图形态、窗口映射机制及其底层实现原理。读完本文,你将能够在流式事件丰富、动态过滤规则、查找表关联等场景中正确使用 side inputs,并理解主/侧输入窗口不一致时的行为。
什么是 Side Inputs
在 Apache Beam 中,ParDo变换处理主输入PCollection中的每个元素时,还可以通过side inputs访问额外的数据。这些附加输入与主输入一起被提供给DoFn,供其在处理过程中读取(见 29_advanced_side_inputs.md)。
Side inputs 的典型价值在于:当管道需要在运行时动态摄取附加数据、而非依赖预设或硬编码值时,它可以基于主PCollection的数据、甚至管道中另一个分支的数据来确定附加数据。最典型的场景是流式分析中的事件丰富(enrichment):用一张查找表(lookup table)为实时到达的流式事件补充维度信息。
从概念上讲,side inputs 与以下机制不同:
- 广播变量 / 配置常量:side inputs 的数据来自管道内的
PCollection(可以是另一个分支、外部数据源读入的结果),而非作业启动前写死的值; - CoGroupByKey 等按 key 的 join:side inputs 不要求主/侧数据共享 key,且每个元素可以独立地、按需地查询整份侧数据(视图)。
Python SDK:把侧输入作为 DoFn 的额外参数
在 Apache Beam Python SDK 中,side inputs 以DoFn.process方法的额外参数或Map/FlatMap变换的额外参数形式传入。Python SDK 支持可选参数、位置参数和关键字参数三种方式(见 29_advanced_side_inputs.md)。
基础形态:
class MyDoFn(beam.DoFn): def process(self, element, side_input): ...五种视图形态:AsSingleton / AsIter / AsList / AsDict / AsMultiMap
传入的参数需要用beam.pvalue下的标记类包装,以声明侧输入以何种形态呈现。这些类定义在仓库 sdks/python/apache_beam/pvalue.py 中(AsSingleton见 L475,AsIter见 L524,AsList见 L555,AsDict见 L578,AsMultiMap见 L602):
| 包装类 | 侧输入呈现形态 | 约束与说明 |
|---|---|---|
AsSingleton(pcoll, default_value=...) | 单个普通值 | 每个窗口必须恰好一个元素;为空时返回默认值或EmptySideInput;多于一个元素会抛ValueError |
AsIter(pcoll) | 可迭代对象(迭代器) | 顺序访问,内存效率高 |
AsList(pcoll) | 列表 | 强制物化为 list,适合需要随机访问或取长度的场景 |
AsDict(pcoll) | 字典(key -> value) | 输入须为 (key, value) 二元组且 key 唯一 |
AsMultiMap(pcoll) | 字典(key -> 值列表) | 允许一个 key 对应多个值,且按需惰性读取;要求输入是 KV 对 |
例如在Map中使用AsSingleton,从仓库示例 map_side_inputs_singleton.py 可以看到完整可运行代码:
import apache_beam as beam with beam.Pipeline() as pipeline: chars = pipeline | 'Create chars' >> beam.Create(['# \n']) plants = ( pipeline | 'Gardening plants' >> beam.Create([ '# 🍓Strawberry\n', '# 🥕Carrot\n', '# 🍆Eggplant\n', '# 🍅Tomato\n', '# 🥔Potato\n', ]) | 'Strip header' >> beam.Map( lambda text, chars: text.strip(chars), chars=beam.pvalue.AsSingleton(chars), ) | beam.Map(print))这里把主输入plants与侧输入chars同时传给Map:chars=beam.pvalue.AsSingleton(chars)以关键字参数形式声明侧输入,Lambda 的首个位置参数仍是主输入元素。
CombineGlobally同样支持侧输入。仓库示例 combineglobally_side_inputs_iter.py 展示了用beam.pvalue.AsIter(exclude)把一份"例外项"列表注入全局聚合逻辑:
| 'Get common items with exceptions' >> beam.CombineGlobally( lambda items, exclude: set(items).difference(*exclude), exclude=beam.pvalue.AsIter(exclude))实战示例:BigQuery 数据作为侧输入
仓库中的 bigquery_side_input.py 演示了更贴近真实业务的做法——把 BigQuery 读出的数据以三种不同形态注入变换:
from apache_beam.pvalue import AsList from apache_beam.pvalue import AsSingleton # 从 BigQuery 读入的 PCollection 作为侧输入 ... | beam.Map( attach_corpus_fn, AsList(corpus), AsSingleton(ignore_corpus)) ... | beam.Map( attach_word_fn, AsList(word), AsSingleton(ignore_word)))该示例从publicdata:samples.shakespeare中随机选取语料与单词生成分组数据,AsList用于需要随机下标访问的语料集合,AsSingleton用于"需忽略的语料/单词"这类单值配置。这印证了 side inputs 的定位:数据可以来自管道内任意分支,包括外部系统读入的结果。
Java SDK:withSideInputs 与 ProcessContext.sideInput
在 Java SDK 中,side inputs 通过ParDo的.withSideInputs(...)方法附加,在DoFn内通过DoFn.ProcessContext.sideInput(view)读取(见 29_advanced_side_inputs.md)。
PCollection<Integer> input = ...; PCollectionView<Integer> sideInput = ...; PCollection<Integer> output = input.apply(ParDo.of(new DoFn<Integer, Integer>() { @ProcessElement public void processElement(ProcessContext c) { Integer sideInputValue = c.sideInput(sideInput); ... } }).withSideInputs(sideInput));withSideInputs的多个重载定义在 ParDo.java 中(单输入形态见 L735-L776,多输出形态见 L900-L938),支持:
- 可变参数
withSideInputs(PCollectionView<?>... sideInputs); Iterable<PCollectionView<?>>;Map<String, PCollectionView<?>>(带 tag 的命名侧输入,便于多输出场景按 tag 区分)。
视图的创建:View 变换族
PCollectionView是"把PCollection呈现为类型T的不可变视图,可作为ParDo的 side input 访问"的接口(见 PCollectionView.java 的类注释 L30-L47)。最常见的是用View变换族来制备视图,其定义在 View.java:
| View 变换 | 侧输入形态 | 约束(源自源码注释) |
|---|---|---|
View.asSingleton() | 单值 | 输入为空时在消费方DoFn抛NoSuchElementException;多于一个元素抛IllegalArgumentException(L157-L161) |
View.asList() | List<T> | 需要随机访问或取大小时使用;顺序访问用asIterable性能更好;部分 Runner 要求视图可装入内存(L167-L176) |
View.asIterable() | Iterable<T> | 顺序访问;部分 Runner 要求装入内存(L181-L186) |
View.asMap() | Map<K, V> | 要求每个窗口内每个 key 唯一;不唯一时先Combine.perKey或改用asMultimap(L192-L205) |
View.asMultimap() | Map<K, Iterable<V>> | 不要求 key 唯一,允许一 key 多值(L213-L225) |
单值视图的典型制备方式(源码注释示例):
PCollection<InputT> input = ...; PCollectionView<OutputT> output = input .apply(Combine.globally(yourCombineFn)) .apply(View.<OutputT>asSingleton());Map 视图通常配合Combine.perKey使用,以保证每个 key 在窗口内只有一个值:
PCollection<KV<K, V>> input = ...; PCollectionView<Map<K, OutputT>> output = input .apply(Combine.perKey(yourCombineFn)) .apply(View.<K, OutputT>asMap());窗口化数据中的 Side Inputs:主窗口投影到侧输入窗口
Side inputs 同样适用于窗口化数据。Apache Beam 使用主输入元素的窗口去查找侧输入元素的对应窗口:将主输入的窗口投影到侧输入的窗口集合上,再从投影得到的窗口中取侧输入值。主输入与侧输入可以拥有相同或不同的窗口策略(见 29_advanced_side_inputs.md)。
例如:主输入PCollection按10 分钟开窗,侧输入按1 小时开窗。处理某个主输入元素时,Beam 会把该元素所属的 10 分钟窗口投影到小时窗口集合,选中包含它的那个 1 小时窗口,读取该小时的侧输入值。这样,同一小时内到达的多个 10 分钟窗口元素,都能共享同一份小时级侧输入数据——非常适合"小时级更新的配置表 + 分钟级事件流"的组合。
窗口映射的底层实现
窗口映射函数由 sideinputs.py 中的default_window_mapping_fn(L53-L66)生成,其行为在源码层面可以确认:
- 若侧输入的窗口函数是
GlobalWindows(),则映射函数固定返回GlobalWindow(_global_window_mapping_fn,L48-L50),即全局窗口侧输入对任意主输入窗口都可见; - 若侧输入使用了
Sessions(会话窗口)则直接抛出RuntimeError,因为无法从任意主输入窗口唯一确定一个会话窗口(L58-L59); - 对于其他窗口类型,映射函数
map_via_end取源窗口的max_timestamp()作为时间点,调用侧输入窗口函数assign重新分配,并取分配结果的最后一个窗口(L61-L64)——这就是"主输入窗口投影到侧输入窗口集合"这一规则的实现。
每个侧输入视图都携带一个WindowMappingFn:AsSideInput.__init__在创建时即调用default_window_mapping_fn(pcoll.windowing.windowfn)绑定该侧输入自己的窗口策略(见 pvalue.py L352-L356)。运行时通过SideInputMap维护"窗口 -> 侧输入值"的映射,读取时先用主输入窗口调用映射函数定位侧输入窗口,再取值(sideinputs.py L77-L84)。
在 Java 侧,PCollectionView.getWindowMappingFn()同样是视图的必备组成部分(PCollectionView.java L81-L88),由各 Runner 在执行时用于窗口投影。
全局窗口与窗口化侧输入的使用建议
- 全局窗口的侧输入:由于映射固定返回
GlobalWindow,它天然对每个窗口的主输入可见,是最省心的组合; - 不同窗口策略的组合:只要侧输入窗口策略可被确定性投影(固定窗口、滑动窗口均可),主/侧输入窗口大小不一致是允许的,但请牢记数据延迟取决于侧输入窗口的触发频率;
- 会话窗口做侧输入会直接报错:Python 实现中会在构建窗口映射函数时抛出
RuntimeError,务必避免。
运行原理:从 PCollection 到可查询的视图
把一条PCollection变成可查询的 side input,在 Python SDK 中经历了清晰的抽象分层(见 pvalue.py):
AsSideInput标记:所有As*包装类的基类,声明"该 PCollection 将作为侧输入使用",并在构造时确定window_mapping_fn与_windowed_coder(L342-L369);SideInputData:封装 side input 的完整规格——访问模式 URN(ITERABLE/MULTIMAP)、窗口映射函数、视图构建函数view_fn,并负责与 Runner API 的 proto 互转(to_runner_api/from_runner_api,L439-L472);- 运行时物化:每种视图实现
_from_runtime_iterable,把运行时得到的可迭代数据转换成实际形态。例如AsSingleton只取前两个元素做校验——为空时返回默认值或EmptySideInput,恰好一个则返回该值,两个及以上抛ValueError(L507-L517)。
值得注意的边界对象EmptySideInput(L630-L639):当单例侧输入所在窗口为空且未提供默认值时,它会被作为占位值传给DoFn。业务代码如需区分"空侧输入"与"真实值",应显式检查该对象。
实践注意事项与最佳实践
综合原文档与仓库实现,使用 side inputs 时请重点关注以下几点:
- 单值视图的元素数约束:
AsSingleton/View.asSingleton()要求每个窗口内恰好一个元素。Python 侧为空时可用default_value兜底;Java 侧空窗口会抛NoSuchElementException,多元素会抛IllegalArgumentException。准备数据时务必先Combine.globally或去重。 - Map 视图的 key 唯一性:
AsDict/View.asMap()要求 key 在窗口内唯一;数据可能重复时,Python 用AsMultiMap,Java 用View.asMultimap(),或先做Combine.perKey。 - 内存与访问方式权衡:多个 Runner 要求视图装入内存。顺序访问优先
AsIter/asIterable;需要随机访问或取长度才用AsList/asList。AsMultiMap采用惰性按 key 读取,适合大数据量查找。 - 会话窗口禁止用作侧输入:Python 实现会直接拒绝
Sessions窗口的侧输入;主输入与侧输入窗口策略不同时,请先确认目标 Runner 对窗口投影的支持。 - 动态性优于硬编码:side inputs 的数据在管道运行时由其他分支或外部 IO(如 BigQuery,见 bigquery_side_input.py)产生,因而可在不重启作业的情况下按窗口粒度刷新,这是流式事件丰富场景的关键优势。
总结
Side inputs 是 Apache Beam 在"主输入 + 附加数据"维度上的核心抽象:Python SDK 以DoFn.process/Map/FlatMap的额外参数配合AsSingleton、AsIter、AsList、AsDict、AsMultiMap五类包装声明侧输入;Java SDK 则以View变换族制备PCollectionView,经ParDo.withSideInputs注入、ProcessContext.sideInput读取。在处理窗口化数据时,Beam 通过窗口映射函数把主输入窗口投影到侧输入窗口集合,从而支持主/侧输入采用不同窗口策略(如 10 分钟主窗口 + 1 小时侧窗口),但需规避会话窗口。理解这些机制,即可在流式事件丰富、动态过滤与查找表关联等场景中写出既正确又高效的 Beam 管道。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Side Inputs 完全指南:给 ParDo 注入运行时数据(Java / Python / Go)
Apache Beam Side Inputs 完全指南:给 ParDo 注入运行时数据(Java / Python / Go) 本文围绕 Tour of Be
大数据批处理流处理数据工程Apache Beam Side Inputs 详解:为 ParDo 注入动态附加数据
Apache Beam Side Inputs 详解:为 ParDo 注入动态附加数据 Apache Beam 的 Side Inputs(侧输入)机制允许 P
批处理流处理大数据Apache Beam Python SDK Side Input 实战:用 ParDo 侧输入为数据流动态注入附加数据
Apache Beam Python SDK Side Input 实战:用 ParDo 侧输入为数据流动态注入附加数据 本篇技术指南聚焦 Apache Bea
批处理流处理大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考