ARTICLE DETAIL

资讯详情

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

Apache Beam 核心变换系列:用简单函数实现 Combine 聚合变换

Apache Beam 核心变换系列:用简单函数实现 Combine 聚合变换 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Combine是 Apache Beam 中用于将PCollection中的一组元素或数值合并为单一结果的变换transform它既有作用于整个PCollection的全局变体如CombineGlobally/Combine.globally/beam.Combine也有针对key/value键值对PCollection的按 key 变体CombinePerKey/Combine.perKey。本文以 simple-function/description.md 为主线讲解在 Java、Python、Go 三种 SDK 下如何用简单函数完成求和等基础聚合并结合仓库源码剖析其底层机制与可替换的内置组合函数。Combine 是什么Combine是 Beam 中用于把数据里的元素或数值合并起来的变换。它的适用场景非常广泛计算一个PCollection中所有整数的总和、最小值、最大值把一组字符串拼接成一个字符串对PCollection中的key/value键值对按 key 分组再把同一 key 下的所有 value 合并此时使用CombinePerKey变体。当你应用一个Combine变换时必须提供一个包含合并逻辑的函数。Beam SDK 同时提供了若干预置的合并函数pre-built combine functions覆盖 sum、min、max 等常见数值运算开箱即用无需自己手写。合并函数必须满足的两个性质文档中特别强调了一个关键约束合并函数必须是可交换commutative且可结合associative的。原因在于该函数不保证对某个 key 的所有 value 恰好只调用一次输入数据包括 value 集合可能被分布到多个 worker 上因此合并函数可能被多次调用以对 value 集合的各个子集执行部分合并partial combining最后再把部分结果汇总。正是这两个数学性质保证了无论数据在分布式环境下如何切分、合并顺序如何变化最终结果都一致。如果合并逻辑依赖元素顺序例如先出现的元素优先就无法安全地用于分布式场景。简单函数求和——三种 SDK 的最小实现对于sum这类简单合并操作通常可以用一个**简单函数simple function**直接实现无需定义复杂的累加器类。下面分别给出 Java、Python、Go 三种 SDK 的完整写法。Gobeam.Combine 内联函数Go SDK 中可以直接传入一个内联的二元函数函数签名形如func(sum, elem int) int第一个参数是当前累计值第二个参数是待合并的元素func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.Combine(s, func(sum, elem int) int { return sum elem }, input) }在仓库对应的可运行示例 main.go 中完整流程是input : beam.Create(s, 10, 30, 50, 70, 90) output : applyTransform(s, input) debug.Print(s, output) err : beamx.Run(ctx, p)输入是10, 30, 50, 70, 90beam.Combine将其合并为总和250并通过debug.Print输出结果。底层 API 定义在 sdks/go/pkg/beam/combine.gobeam.Combine(s, combinefn, col, opts...)对整列元素做全局合并beam.CombinePerKey(s, combinefn, col, opts...)则按键合并二者均以任意函数作为合并逻辑。JavaCombine.globallySerializableFunctionJava SDK 中简单函数需要实现SerializableFunctionIterableInteger, Integer接口输入是元素的迭代集合输出是合并后的单一结果。文档中的SumInts即为标准写法// Sum a collection of Integer values. The function SumInts implements the interface SerializableFunction. public static class SumInts implements SerializableFunctionIterableInteger, Integer { Override public Integer apply(IterableInteger input) { int sum 0; for (int item : input) { sum item; } return sum; } }使用方式取自 Task.javaPCollectionInteger input pipeline.apply(Create.of(10, 30, 50, 70, 90)); PCollectionInteger output input.apply(Combine.globally(new SumIntegerFn()));Combine.globally是 Java 侧全局合并的入口。从 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Combine.java 可以看到它提供了多种重载接收SerializableFunctionIterableV, V、SerializableBiFunctionV, V, V二元函数即 BinaryCombineFn 风格以及带CombineFn的通用版本Combine.perKey则对应按 key 合并同文件 L165-L200。Pythonbeam.CombineGlobally 普通函数Python SDK 最为灵活直接传入普通 Python 函数即可并且支持带默认参数的函数。文档示例通过bound参数演示了这一点input [1, 10, 100, 1000] def bounded_sum(values, bound500): return min(sum(values), bound) small_sum input | beam.CombineGlobally(bounded_sum) # [500] large_sum input | beam.CombineGlobally(bounded_sum, bound5000) # [1111]small_sum使用默认bound50011010010001111被截断为500因此输出[500]large_sum传入bound5000总和1111未超限原样输出[1111]。beam.CombineGlobally是 Python SDK 的全局合并变换定义于 sdks/python/apache_beam/transforms/core.py。它支持丰富的调用方式CombineGlobally(sum)直接使用函数、CombineGlobally(sum).as_singleton_view()把结果作为单例侧输入side input使用、CombineGlobally(sum).without_defaults()在输入为空时不输出默认值避免空 PCollection 被填充默认结果参考同文件 L3020-L3033 的相关说明。仓库中对应的可运行示例 task.py 采用同样模式with beam.Pipeline() as p: (p | beam.Create([1, 2, 3, 4, 5]) | beam.CombineGlobally(sum) | Output())不止数字Combine 同样适用于字符串等任意类型文档强调Combine的输入数据由整数构成但它也可以与其他类型结合使用例如strings及其他类型。下面的练习示例把 8 个英文单词合并成一个逗号分隔的字符串。Go 版本input : beam.Create(s, quick, brown, fox, jumps, over, the, lazy, dog) func applyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.Combine(s, func(sum, elem string) string { return sum,elem }, input) }Java 版本同样只需把SerializableFunction的泛型换成String在累加循环中用StringBuilder拼接public static class ConcatenateStrings implements SerializableFunctionIterableString, String { Override public String apply(IterableString input) { StringBuilder concatenated new StringBuilder(); for (String item : input) { concatenated.append(,).append(item); } return concatenated.toString(); } } PCollectionString input pipeline.apply(Create.of(quick, brown, fox, jumps, over, the, lazy, dog)); PCollectionString concatenated input.apply(Combine.globally(new ConcatenateStrings()));Python 版本def concat(strings): total for word in strings: total word return total with beam.Pipeline() as p: (p | beam.Create([quick, brown, fox, jumps, over, the, lazy, dog]) | beam.CombineGlobally(concat) | LogElements())从这些例子可以看出简单函数模式的本质是接收整个元素的集合返回一个合并结果。只要该逻辑满足可交换、可结合Beam 就能在分布式执行中安全地分批合并最终结果与串行执行完全一致。从简单函数到通用 CombineFn简单函数适合sum、字符串拼接这类输入类型与输出类型相同的场景。一旦合并逻辑变复杂例如计算平均值输出类型是浮点数、中间态是总和 计数就需要创建CombineFn子类通过四个方法定义完整的合并协议accumulation type 可以不同于输入/输出类型Create Accumulator创建新的局部累加器。求平均值的例子中局部累加器记录已累加值的总和最终除法中的分子与已累加值的个数分母可在分布式场景下被任意多次调用Add Input把单个输入元素加入累加器并返回新的累加器示例中更新总和并自增计数同样可以并行调用Merge Accumulators把多个累加器合并为一个平均值场景中合并各部分的分子与分母其输出还可能被再次合并任意多次Extract Output执行最终计算平均值即总和除以个数只在最终合并后的累加器上调用一次。关于这一主题的完整三 SDK 实现可继续阅读同组的 combine-fn/description.md对每个 key 下的 value 合并以及基于二元函数的BinaryCombineFn写法可分别参考 combine-per-key/description.md 与 binary-combine-fn/description.md。四个单元共同构成 Combine 完整课程见 group-info.yaml。预置合并函数与 Runner 优化对于 sum、min、max 等常见数值操作不必每次手写函数Beam SDK 提供了一系列预置合并函数例如Sum、Min、Max家族Java 的Combine.globally(Sum.integers())、Python 的beam.CombineGlobally(beam.combiners.Sum.IntsFn())、Go 的combine.SumInts()等它们天然满足可交换、可结合约束语义清晰且类型安全。之所以强调可交换、可结合还因为这两条性质允许 Runner 自动应用关键优化详见 combine-fn/description.mdCombiner lifting最显著的优化在数据被 shuffle 之前先在每个 key、每个窗口内完成部分合并从而把需要 shuffle 的数据量减少数个数量级也称mapper-side combineIncremental combining增量合并当CombineFn能显著缩减数据大小时在流式 shuffle 过程中边产出边合并把合并开销摊薄到计算空闲时段同时降低中间累加器的存储占用。理解这些机制有助于你判断手写简单函数时务必遵守交换律与结合律否则 Runner 的任何并行/部分合并优化都可能改变结果语义。小结与延伸阅读简单函数是入门 Combine 的最佳切入点三 SDK 写法Go 用beam.Combine(s, fn, input)内联函数Java 用Combine.globally(new SerializableFunction...)Python 用input | beam.CombineGlobally(fn)还支持带默认参数的函数约束合并函数必须可交换、可结合因为它在分布式环境下可能被多次调用做部分合并通用性不仅支持整数求和字符串拼接等任意类型的合并同样适用进阶复杂合并升级为CombineFn四方法createAccumulator / addInput / mergeAccumulators / extractOutput按 key 合并使用CombinePerKey二元合并可基于BinaryCombineFn内置能力优先考虑 SDK 预置的 sum / min / max 合并函数并在理解 combiner lifting 与增量合并等 Runner 优化后合理设计自己的合并逻辑。本文对应的可运行示例位于 simple-function 目录含 Go 示例、Java 示例、Python 示例欢迎在 Beam Playground 中运行并尝试修改输入数据与合并逻辑验证 Combine 的分布式合并语义。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 核心变换实战用 CombinePerKey 按用户名聚合游戏得分Tour of Beam 挑战二Apache Beam 核心变换实战用 CombinePerKey 按用户名聚合游戏得分Tour of Beam 挑战二 本篇技术指南围绕 Apache大数据批处理流处理数据工程Apache Beam Go SDK Kata 实战用 Combine 简单函数实现求和Apache Beam Go SDK Kata 实战用 Combine 简单函数实现求和 Combine 是 Apache Beam 中用于把集合中的元素或值Apache Beam Kotlin Katas 实战Aggregation 之 Count 聚合变换详解Apache Beam Kotlin Katas 实战Aggregation 之 Count 聚合变换详解 导读 本文以 Apache Beam 官方 Kot大数据批处理流处理数据工程上一篇彻底解决Khoj项目Docker部署中的Bad Request(400)错误从根源排查到完美修复下一篇深入理解Huihui-Qwythos-9B-Claude-Mythos-5-1M-abliterated-GGUF函数调用与工具使用实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表