ARTICLE DETAIL

资讯详情

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

Apache Beam Go SDK 实战:使用 filter 包(Include/Exclude)过滤 PCollection 元素

Apache Beam Go SDK 实战:使用 filter 包(Include/Exclude)过滤 PCollection 元素 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载导读本文聚焦 Apache Beam Go SDK 中 transforms/filter 包的用法承接用 ParDo 过滤元素这一前置课程介绍如何用现成的filter.Include与filter.Exclude变换按任意布尔谓词predicate批量保留或剔除 PCollection 中的元素。读完本文你将掌握 filter 包的函数签名、语义差异、底层执行原理并能独立完成 katas 练习中过滤出奇数/偶数的实战任务。1. 课程背景从 ParDo 手写过滤到专用 filter 包在 common_transforms/filter/pardo 这节课中我们曾用beam.ParDo配合自定义DoFn手动实现过滤逻辑——在ProcessElement里判断条件并决定是否emit元素。这是理解 Beam 基本执行模型的重要一步但每次过滤都要手写 DoFn 略显繁琐。本课 filter/filter/task.md 引入的filter包正是为按条件移除管道元素这一高频需求准备的专用工具。它内部仍然基于beam.ParDo实现但把条件判定 发射/丢弃的样板代码封装成了两个可直接调用的变换函数filter.Include(s, col, fn)保留使fn返回true的元素filter.Exclude(s, col, fn)剔除使fn返回true的元素与 Include 语义相反。从课程编排 lesson-info.yaml 可以看到pardo与filter是本 lesson 下的两个递进小节先手动、后封装形成完整的知识闭环。2. Kata 任务定义本课任务是Kata实现一个过滤函数使用filter包中的方法把奇数odd numbers过滤掉。提示hint指出使用filter.Exclude。任务骨架文件 pkg/task/task.go 中函数体位置以TODO()占位对应 task-info.yaml 中定义的 placeholder需要你补全ApplyTransform的实现func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return filter.Exclude(s, input, func(element int) bool { return element % 2 1 }) }关键点解读ApplyTransform的签名是(beam.Scope, beam.PCollection) beam.PCollection即输入一个 PCollection输出过滤后的新 PCollection谓词函数接收元素element int返回boolelement % 2 1判断元素是否为奇数而Exclude会剔除谓词返回true的元素因此奇数被过滤掉偶数被保留——这正是过滤掉奇数的含义。如果改用filter.Include则语义翻转Include保留谓词为true的元素此时需要把谓词改为element%2 0才能达到同样效果。两种写法结果等价选择哪一个取决于你更习惯表达要什么还是不要什么。3. 完整可运行的 Pipeline 示例Kata 的可执行入口 cmd/main.go 展示了 filter 变换在完整管道中的用法package main import ( context beam.apache.org/learning/katas/common_transforms/filter/filter/pkg/task github.com/apache/beam/sdks/v2/go/pkg/beam github.com/apache/beam/sdks/v2/go/pkg/beam/log github.com/apache/beam/sdks/v2/go/pkg/beam/x/beamx github.com/apache/beam/sdks/v2/go/pkg/beam/x/debug ) func main() { ctx : context.Background() p, s : beam.NewPipelineWithRoot() input : beam.Create(s, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10) output : task.ApplyTransform(s, input) debug.Print(s, output) err : beamx.Run(ctx, p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }运行流程解析beam.NewPipelineWithRoot()创建管道对象p与根作用域sbeam.Create(s, 1, 2, 3, ..., 10)构造一个包含 1~10 十个整数的输入 PCollectiontask.ApplyTransform(s, input)执行你实现的 filter 变换debug.Print(s, output)将结果打印到日志等价于在本地直接观察输出beamx.Run(ctx, p)在当前环境如本地 Direct runner执行整个管道。运行后预期输出为过滤掉奇数后的偶数序列2, 4, 6, 8, 10。4. 单元测试验证passert.Equals 断言结果Kata 的测试文件 test/task_test.go 对实现做了精确校验func TestApplyTransform(t *testing.T) { p, s : beam.NewPipelineWithRoot() tests : []struct { input beam.PCollection want []interface{} }{ { input: beam.Create(s, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), want: []interface{}{2, 4, 6, 8, 10}, }, } for _, tt : range tests { got : task.ApplyTransform(s, tt.input) passert.Equals(s, got, tt.want...) if err : ptest.Run(p); err ! nil { t.Error(err) } } }测试要点使用beam.Create构造输入将期望结果写成[]interface{}{2, 4, 6, 8, 10}passert.Equals(s, got, tt.want...)是 Beam Go SDK 的断言工具它声明got必须等于want在管道执行阶段进行校验ptest.Run(p)在测试环境中运行管道。任何元素缺失或多余都会导致断言失败从而保证 filter 逻辑的精确性。这正是测试驱动的 kata 练习方式只要你的ApplyTransform实现正确测试即通过。5. 源码级原理Include / Exclude 是如何工作的要深入理解 filter 包直接阅读其实现 sdks/go/pkg/beam/transforms/filter/filter.go 是最佳途径。5.1 谓词签名约束filter 包在初始化时定义了统一的谓词签名sig funcx.MakePredicate(beam.TType) // T - bool即过滤函数必须是一元函数接收一个类型为T的元素返回bool。Include和Exclude内部都会调用funcx.MustSatisfy(fn, ...)在管道构建期强制校验这个签名如果传入的函数形态不符会立即报错而非等到运行时。5.2 Include 与 Exclude 的差异两者的实现几乎一致唯一的区别是一个布尔标记func Include(s beam.Scope, col beam.PCollection, fn any) beam.PCollection { s s.Scope(filter.Include) funcx.MustSatisfy(fn, funcx.Replace(sig, beam.TType, col.Type().Type())) return beam.ParDo(s, filterFn{Predicate: beam.EncodedFunc{Fn: reflectx.MakeFunc(fn)}, Include: true}, col) } func Exclude(s beam.Scope, col beam.PCollection, fn any) beam.PCollection { s s.Scope(filter.Exclude) funcx.MustSatisfy(fn, funcx.Replace(sig, beam.TType, col.Type().Type())) return beam.ParDo(s, filterFn{Predicate: beam.EncodedFunc{Fn: reflectx.MakeFunc(fn)}, Include: false}, col) }注意这里二者都先创建子作用域filter.Include/filter.Exclude便于在管道图中形成清晰的变换节点都通过reflectx.MakeFunc(fn)把 Go 函数编码为beam.EncodedFunc使其可以随管道描述被序列化、跨 runner 传递返回的 PCollection 与输入类型一致过滤不改变元素类型。5.3 filterFn 的执行逻辑核心执行逻辑封装在filterFnDoFn 中type filterFn struct { Predicate beam.EncodedFunc json:predicate Include bool json:include fn reflectx.Func1x1 } func (f *filterFn) Setup() { f.fn reflectx.ToFunc1x1(f.Predicate.Fn) } func (f *filterFn) ProcessElement(elm beam.T, emit func(beam.T)) { match : f.fn.Call1x1(elm).(bool) if match f.Include { emit(elm) } }这段代码揭示了关键设计Setup()阶段把编码后的谓词函数还原为可调用的Func1x1反序列化时机与 runner 的生命周期管理有关ProcessElement对每个元素调用谓词得到match判定规则if match f.IncludeIncludetrue时谓词为true才发射Includefalse即 Exclude时谓词为false才发射。也就是说Exclude 的本质是谓词取反后 Include两者共用一个 DoFn 实现只是开关不同。另外包初始化时通过register机制注册了filterFn及关联的发射器register.DoFn2x0[beam.T, func(beam.T)]、register.Emitter1[beam.T]()等这是 Beam Go SDK 为了保证函数可序列化、可在分布式 runner如 Flink、Dataflow、Spark 等上执行所做的注册要求。5.4 关于谓词函数形态的注意事项filter.go 的文档注释明确提醒Filter functions must be registered with Beam, and must not be closures.即过滤函数需要注册且不能是闭包closure官方示例采用命名函数 register.Function1x1的写法func lessThanThree(s string) bool { return len(s) 3 } // 注册过滤函数 func init() { register.Function1x1(lessThanThree) } words : beam.Create(s, a, b, long, alsolong) short : filter.Include(s, words, lessThanThree) // short 在运行时包含 a 和 b这一约束的背景是Beam 管道描述需要跨进程/跨机器传输匿名闭包捕获的外部变量无法被可靠序列化。而 kata 练习中采用内联匿名函数func(element int) bool { return element % 2 1 }在本地 Direct runner 场景下可以正常工作适合教学演示在构建生产级、需要跨 runner 移植的管道时请遵循注册命名函数的最佳实践。6. Include 与 Exclude 的语义对照速查变换谓词返回 true谓词返回 false典型用法filter.Include(s, col, fn)保留剔除只要满足条件的元素如len(s) 3的短字符串filter.Exclude(s, col, fn)剔除保留排除满足条件的元素如本课的过滤掉奇数两者的返回 PCollection 元素类型与输入一致谓词签名固定为T - bool。选 Include 还是 Exclude本质上取决于你的谓词天然表达的是保留条件还是排除条件选择表达更自然的一方即可不必刻意强求。7. 如何运行与验证本课 Kata在仓库根目录的 Go 模块go.mod 与 learning/katas/go/go.mod环境中可以按以下方式练习和验证打开 pkg/task/task.go在ApplyTransform中补全filter.Exclude实现运行入口程序观察输出是否为偶数序列go run ./common_transforms/filter/filter/cmd运行单元测试验证正确性go test ./common_transforms/filter/filter/test测试通过的标准输入1~10输出必须恰好等于2, 4, 6, 8, 10。你也可以自行修改输入集合与谓词例如改成过滤偶数、过滤大于 5 的数体会 Include / Exclude 的语义差异。8. 小结本课的核心收获可归纳为三点API 层面filter.Include/filter.Exclude是 Beam Go SDK 内置的过滤变换签名统一为(Scope, PCollection, T - bool) - PCollection返回类型与输入一致语义层面Include 保留满足谓词的元素Exclude 剔除满足谓词的元素二者共用同一 DoFn 实现见 filter.go仅通过Include布尔标记区分行为工程层面谓词函数需遵循签名T - bool、建议命名并注册、避免闭包的约束才能保证管道可序列化、可在各类 runner 上稳定执行。掌握了 filter 包你就拥有了比手写 ParDo 更简洁、更不易出错的元素筛选手段。接下来可以继续挑战 Go katas 中关于其他 common transforms 的课程把 Beam 的变换工具箱逐步补齐。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Go SDK 实战使用 filter 包从 PCollection 中过滤元素Kata 深度解析Apache Beam Go SDK 实战使用 filter 包从 PCollection 中过滤元素Kata 深度解析 本篇技术指南围绕 Beam Go大数据批处理流处理数据工程Apache Beam Go SDK 聚合实战用 stats.CountElms 统计 PCollection 元素总数Apache Beam Go SDK 聚合实战用 stats.CountElms 统计 PCollection 元素总数 本指南基于 Apache Beam大数据批处理流处理数据工程Apache Beam Go SDK 实战使用 stats.Sum 计算 PCollection 元素总和Kata 详解Apache Beam Go SDK 实战使用 stats.Sum 计算 PCollection 元素总和Kata 详解 本文围绕 Beam Katas大数据批处理流处理数据工程上一篇3个常见性能陷阱与突破方案打造流畅的微信小程序数据可视化下一篇GenWealth 应用接入 AlloyDB Omni在 GKE 上部署本地可移植版 AlloyDB 并配置 AlloyDB AI 向量检索全流程指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表