ARTICLE DETAIL

资讯详情

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

Apache Beam Side Inputs 完全指南:给 ParDo 注入运行时数据(Java / Python / Go)

Apache Beam Side Inputs 完全指南:给 ParDo 注入运行时数据(Java / Python / Go) 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本文围绕 Tour of Beam 核心变换课程中的 side-inputs 单元 展开。Side input侧输入是 Apache Beam 中除主输入PCollection之外为ParDo提供的附加数据源它允许每个DoFn在处理主输入元素时读取一份视图化的辅助数据。读完本文你将掌握 Java / Python / Go 三种 SDK 下侧输入的完整用法、侧输入与窗口windowing的投影规则并能在 Playground 中动手实践城市—国家映射与城市—GMT 时区两类经典练习。一、什么是 Side Input主输入之外的动态数据注入在 Beam 中ParDo变换通常只有一个主输入PCollectionDoFn逐元素处理该集合。但很多业务场景要求在逐元素处理时额外注入一份在运行时才确定的数据——例如一个阈值、一份映射表、一段参考数据。如果硬编码进代码就无法根据输入数据或管线的不同分支动态变化。Side input 正是为这一需求而生它是一个附加输入DoFn在处理主输入PCollection中每个元素时都可以访问指定 side input 时你会把某份数据视图化create a view该视图可在ParDo的DoFn内被读取适用场景附加数据需要在运行时确定而非硬编码例如由输入数据本身推导或依赖管线中另一个分支的输出。从 Java SDK 源码看Beam 对视图的定义非常精确View.java 的类注释写道虽然PCollectionElemT每个窗口有多个ElemT值但PCollectionViewViewT每个窗口只有一个ViewT值可以把它理解成从窗口到ViewT值的映射。当ParDo处理主输入中处于窗口w的元素并通过ProcessContext#sideInput读取视图时返回的正是该窗口w对应的视图值。二、把 Side Input 传给 ParDo三种 SDK 的完整写法2.1 JavawithSideInputsProcessContext.sideInputJava 的核心链路分三步先用Combine.globally聚合出单值再调用.asSingletonView()创建PCollectionView最后在ParDo上用.withSideInputs(...)注入并在DoFn内通过ProcessContext.sideInput(view)读取// Pass side inputs to your ParDo transform by invoking .withSideInputs. // Inside your DoFn, access the side input by using the method DoFn.ProcessContext.sideInput. // The input PCollection to ParDo. PCollectionString words ...; // A PCollection of word lengths that well combine into a single value. PCollectionInteger wordLengths ...; // Singleton PCollection // Create a singleton PCollectionView from wordLengths using Combine.globally and View.asSingleton. final PCollectionViewInteger maxWordLengthCutOffView wordLengths.apply(Combine.globally(new Max.MaxIntFn()).asSingletonView()); // Apply a ParDo that takes maxWordLengthCutOffView as a side input. PCollectionString wordsBelowCutOff words.apply(ParDo .of(new DoFnString, String() { ProcessElement public void processElement(Element String word, OutputReceiverString out, ProcessContext c) { // In our DoFn, access the side input. int lengthCutOff c.sideInput(maxWordLengthCutOffView); if (word.length() lengthCutOff) { out.output(word); } } }).withSideInputs(maxWordLengthCutOffView) );关键点withSideInputs可接受多个视图sideInput的返回类型由视图类型决定——单值视图返回T列表视图返回ListT映射视图返回MapK, V详见下文视图的几种形态。2.2 Python额外参数 pvalue.As*延迟展开Python SDK 的侧输入以DoFn.process方法或Map/FlatMap的可调用对象的额外参数形式传入。可选参数、位置参数、关键字参数全部支持构造期传入的延迟参数deferred arguments会被解包为实际值。例如pvalue.AsIter(pcoll)会在每次process调用时传入pcoll实际元素的可迭代对象# Side inputs are available as extra arguments in the DoFns process method or Map / FlatMaps callable. # Optional, positional, and keyword arguments are all supported. Deferred arguments are unwrapped into their # actual values. For example, using pvalue.AsIter(pcoll) at pipeline construction time results in an iterable # of the actual elements of pcoll being passed into each process invocation. In this example, side inputs are # passed to a FlatMap transform as extra arguments and consumed by filter_using_length. words ... # Callable takes additional arguments. def filter_using_length(word, lower_bound, upper_boundfloat(inf)): if lower_bound len(word) upper_bound: yield word # Construct a deferred side input. avg_word_len ( words | beam.Map(len) | beam.CombineGlobally(beam.combiners.MeanCombineFn())) # Call with explicit side inputs. small_words words | small beam.FlatMap(filter_using_length, 0, 3) # A single deferred side input. larger_than_average ( words | large beam.FlatMap( filter_using_length, lower_boundpvalue.AsSingleton(avg_word_len)) ) # Mix and match. small_but_nontrivial words | beam.FlatMap( filter_using_length, lower_bound2, upper_boundpvalue.AsSingleton(avg_word_len)) # We can also pass side inputs to a ParDo transform, which will get passed to its process method. # The first two arguments for the process method would be self and element. class FilterUsingLength(beam.DoFn): def process(self, element, lower_bound, upper_boundfloat(inf)): if lower_bound len(element) upper_bound: yield element small_words words | beam.ParDo(FilterUsingLength(), 0, 3)上例演示了三种典型用法普通字面量参数0, 3、单个延迟侧输入pvalue.AsSingleton(avg_word_len)作为关键字参数、以及字面量与侧输入混用。Python 的解包机制在源码中有清晰对应pvalue.py 导出了AsSingleton、AsIter、AsList、AsDict、AsMultiMap五类侧输入封装它们都继承自AsSideInput。2.3 Gobeam.SideInput 迭代器签名的 DoFnGo SDK 中侧输入通过beam.SideInput{Input: ...}传入beam.ParDo在DoFn的ProcessElement方法中以迭代器参数形式出现排在主输入元素之后、输出发射器之前。此外所有侧输入可迭代对象都应使用泛型register.IterX[...]注册以优化运行时执行// Side inputs are provided using beam.SideInput in the DoFns ProcessElement method. // Side inputs can be arbitrary PCollections, which can then be iterated over per element // in a DoFn. // Side input parameters appear after main input elements, and before any output emitters. words ... // avgWordLength is a PCollection containing a single element, a singleton. avgWordLength : stats.Mean(s, wordLengths) // Side inputs are added as with the beam.SideInput option to beam.ParDo. wordsAboveCutOff : beam.ParDo(s, filterWordsAbove, words, beam.SideInput{Input: avgWordLength}) wordsBelowCutOff : beam.ParDo(s, filterWordsBelow, words, beam.SideInput{Input: avgWordLength}) // filterWordsAbove is a DoFn that takes in a word, // and a singleton side input iterator as of a length cut off // and only emits words that are beneath that cut off. // // If the iterator has no elements, an error is returned, aborting processing. func filterWordsAbove(word string, lengthCutOffIter func(*float64) bool, emitAboveCutoff func(string)) error { var cutOff float64 ok : lengthCutOffIter(cutOff) if !ok { return fmt.Errorf(no length cutoff provided) } if float64(len(word)) cutOff { emitAboveCutoff(word) } return nil } // filterWordsBelow is a DoFn that takes in a word, // and a singleton side input of a length cut off // and only emits words that are beneath that cut off. // // If the side input isnt a singleton, a runtime panic will occur. func filterWordsBelow(word string, lengthCutOff float64, emitBelowCutoff func(string)) { if float64(len(word)) lengthCutOff { emitBelowCutoff(word) } } func init() { register.Function3x1(filterWordsAbove) register.Function3x0(filterWordsBelow) // 1 input of type string Emitter1[string] register.Emitter1[string]() // 1 input of type float64 Iter1[float64] register.Iter1[float64]() } // The Go SDK doesnt support custom ViewFns. // See https://github.com/apache/beam/issues/18602 for details // on how to contribute them!注意 Go 的两个典型约束单值视图既可以按迭代器读取func(*float64) bool空则报错中止也可以直接按值读取float64若非单值会触发运行时 panicGo SDK 目前不支持自定义ViewFn详见上述注释中的 upstream issue。三、视图的几种形态单值、列表、映射与多重映射Beam 支持把每个窗口内的PCollection视图化为不同形态。Java SDK 的 View.java 明确给出了五种视图变换Python 则在pvalue.As*家族中一一对应视图形态Java 变换Python 封装适用场景单值SingletonView.asSingleton()pvalue.AsSingletonCombine.globally聚合出的单个值如阈值、平均值列表ListView.asList()pvalue.AsList窗口可整体装入内存的小型PCollection读取时整个列表缓存可迭代IterableView.asIterable()pvalue.AsIter需要遍历窗口内全部元素映射MapView.asMap()pvalue.AsDictKVK, V且每个键每窗口仅一个值多重映射MultimapView.asMultimap()pvalue.AsMultiMapKVK, V且每个键对应多个值读作MapK, IterableV从源码实现看Java 的View.asMap()与View.asMultimap()正是实现以主输入做查找型 join的推荐方式——当侧输入足够小、能装入内存时用映射视图做键值查找比CoGroupByKey更轻量View.java。四、Side Input 与窗口Windowing的投影规则带窗口的PCollection可能是无限的无法压缩成单一值。当你为带窗口的PCollection创建PCollectionView时该视图表示的是每个窗口一个实体每窗口一个单值、每窗口一个列表……。Beam 使用主输入元素所在的窗口去查找侧输入元素的对应窗口Beam 将主输入元素的窗口投影project到侧输入的窗口集合上取投影结果窗口中的侧输入值若主输入与侧输入窗口完全一致投影得到的就是精确对应的窗口若窗口不同Beam 用投影选择最合适的侧输入窗口。示例主输入使用 1 分钟固定时间窗口侧输入使用 1 小时固定时间窗口则 Beam 会把主输入窗口投影到小时窗口集合从对应的小时窗口取出侧输入值。多窗口与多触发若主输入元素同时存在于多个窗口processElement会对每个窗口各调用一次每次调用都会针对当前窗口投影因此不同窗口调用可能拿到不同的侧输入视图若侧输入有多次触发trigger firingBeam 使用最近一次触发产生的值。这一特性在单全局窗口 自定义 trigger的组合下尤其有用——例如周期性刷新一份参考数据视图。五、Playground 练习从城市—国家到城市—GMT 时区课程的完整可运行代码都在 Playground 窗口里可直接运行和实验。练习的数据结构是入口处有一张以城市为键、国家为值的映射以及一个带name和city字段的Person结构。目标是比较城市并给Person嵌入国家乃至当地时间。三个 SDK 的完整练习实现分别位于仓库中go-example/main.goapplyTransform把citiesToCountries包成beam.SideInput视图joinFn用func(*string, *string) bool迭代器遍历城市映射、命中城市后发射带国家的Personjava-example/Task.javacreateView用View.asMap()把KVString, String视图化为PCollectionViewMapString, StringapplyTransform中context.sideInput(...)查表并out.output新的Personpython-example/task.pyEnrichCountryDoFn直接以cities_to_countries字典作为process的附加参数用cities_to_countries[element.city]完成国家填充。作为延伸实验可以把国家换成GMT 时差再给每个Person附上当地时间Go 版本——先改数据源为时区映射再在joinFn中叠加当前时间citiesToTimeKV : beam.ParDo(s, func(_ []byte, emit func(string, int)){ emit(Beijing, 8) emit(London, 0) emit(San Francisco, -8) emit(Singapore, 8) emit(Sydney, 11) }, beam.Impulse(s)) func joinFn(person Person, citiesToCountriesIter func(*string,*int) bool, emit func(Person)) { var city string var gmt int now : time.Now() for citiesToCountriesIter(city,gmt) { time : now.Hour()gmt if person.City city { if time 0 { time 24 (now.Hour() gmt) } emit(Person{ Name: person.Name, City: city, Time: (fmt.Sprintf(%d:%d,time,now.Minute())), }) break } } }Go 版开始前需先引入依赖fmt和time。Java 版本——把数据源换成KVString, Integer时区对再在ParDo中计算并格式化当地时间PCollectionKVString, Integer citiesToTimeKV pipeline .apply(ParseCitiesToTimeKV, Create.of( KV.of(Beijing, 8), KV.of(London, 0), KV.of(San Francisco, -8), KV.of(Singapore, 8), KV.of(Sydney, 11) )); PCollectionPerson joined persons .apply(Join, Join.innerJoin(citiesToTimeKV)) .apply(FormatTime, ParDo.of(new DoFnKVString, KVPerson, Integer, Person() { ProcessElement public void processElement(ProcessContext c) { Person person c.element().getValue().getKey(); Integer gmt c.element().getValue().getValue(); LocalTime now LocalTime.now(); int time now.getHour() gmt; if(time 0) { time 24 (now.getHour() gmt); } c.output(new Person(person.name, person.city, String.format(%d:%d, time, now.getMinute()))); } }));Python 版本——用Create建立人员与时区两张集合再在join_fn中做合并与时间换算from apache_beam.transforms.util import Create from apache_beam.transforms.core import ParDo from apache_beam.transforms.join import CoGroupByKey from datetime import datetime pipeline beam.Pipeline() persons p | CreatePersons Create([ {name: John, city: Beijing}, {name: Mary, city: Singapore}, {name: Bob, city: Sydney} ]) citiesToTimeKV p | CreateCitiesToTimeKV Create([ (Beijing, 8), (London, 0), (San Francisco, -8), (Singapore, 8), (Sydney, 11) ]) def join_fn(persons, cities_to_time): person, time persons[0], cities_to_time[0][1] now datetime.now() hour now.hour time if hour 0: hour 24 hour return { name: person[name], city: person[city], time: f{hour}:{now.minute} }六、深入理解视图是窗口到值的映射要真正掌握 side input需要理解其底层抽象。Java 的 View.java 注释点明了本质PCollectionViewViewT是从窗口到ViewT值的映射而View系列变换做的就是把窗口内的ElemT值们转换成该窗口的一个ViewT。Python 端对应物是 pvalue.py 中的AsSideInput基类及其五个子类——它们标记了一个PCollection应被如何视图化并延迟解包Go 端则通过register.IterX/EmitterX的注册表把 DoFn 签名与运行时类型绑定见上文init()示例。这意味着两点工程启示窗口一致性决定正确性side input 永远以主输入元素的当前窗口为查找键。若主输入是无限流、侧输入却是全局窗口请记得用 trigger 定期刷新侧输入视图取最近一次触发值内存边界决定形态选择列表 / 映射视图要求窗口数据可装入内存适合小表 join 大流的经典模式而单值视图如Combine.globally后的均值、最大值是最轻量、最常用的形态。延伸阅读本单元三语言完整可运行示例Go、Java、Python课程单元元数据SDK 覆盖与复杂度说明unit-info.yamlJava SDK 视图变换源码View.javaPython SDK 侧输入封装pvalue.py相邻核心变换课程地图map、聚合combine、额外输出additional-outputs、分支branching等均位于 core-transforms 目录赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Python Kata 实战用 Side Input 为 ParDo 注入运行时附加数据Apache Beam Python Kata 实战用 Side Input 为 ParDo 注入运行时附加数据 导读 Side Input侧输入是 Ap大数据批处理流处理数据工程Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时附加数据Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时附加数据 Side Input侧输入是 Apache Be大数据批处理流处理数据工程Apache Beam Kotlin 实战用 Side Input 在 ParDo 中注入运行时附加数据Apache Beam Kotlin 实战用 Side Input 在 ParDo 中注入运行时附加数据 Side Input侧输入是 Apache Be大数据批处理流处理数据工程上一篇6个专业技巧彻底驯服华硕笔记本风扇噪音下一篇ScaleDown压缩率优化技巧提升Token节省效果的10个专业方法创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表