ARTICLE DETAIL

资讯详情

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

Apache Beam WithKeys 实战:用 Python 将 PCollection 元素转换为键值对的 Kata 精讲

Apache Beam WithKeys 实战:用 Python 将 PCollection 元素转换为键值对的 Kata 精讲 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的WithKeys是一个把PCollection中每个元素包装成(K, V)键值对的轻量变换是后续GroupByKey、CombinePerKey等分组聚合操作的前提。本文以仓库中 learning/katas/python/Common Transforms/WithKeys/WithKeys/task.md 这一 Kata 练习为骨架结合官方源码与单元测试完整讲解任务目标、标准解法、验证方式以及WithKeys的底层实现原理读完后你可以在自己的 Beam 管道中熟练运用WithKeys完成各种造键需求。一、Kata 任务背景在训练营中学习 WithKeysApache Beam 仓库的learning/katas目录是一个面向初学者的编程训练营Kata分为 Python、Java、Kotlin、Go 四种语言实现其中 Python 版位于 learning/katas/python。每个小节由task.md题目描述、task.py代码骨架/参考解答、tests/test_task.py自动判题以及task-info.yaml、lesson-info.yaml课程元数据组成借助 PyCharm Education / EduTools 插件即可打开并逐步闯关具体安装步骤见 learning/katas/python/README.md。本节 WithKeys 属于Common Transforms常见变换章节中的基础题目课程元数据 task-info.yaml 将其标记为类别Core Transforms核心变换复杂度BASIC基础标签map、strings从标签可以看出WithKeys本质上是Map系列变换的一种特化它把输入元素作为值Value再为每个元素计算或指定一个键Key从而把普通元素转换成 Beam 管道中最核心的数据形态——KV 对。二、题目要求把水果名变成首字母键值对task.md的题目表述非常精炼Kata:Convert each fruit name into a key/value pair of its first letter and itself, e.g.apple (a, apple)将每个水果名转换为由它的首字母和它自身组成的键值对例如apple (a, apple)题目的提示hint明确指向官方 Python SDK 文档中的apache_beam.transforms.util.WithKeys。要完成该练习需要将[apple, banana, cherry, durian, guava, melon]这 6 个水果名逐一遍历对每个字符串取首字符word[0:1]作为键原字符串作为值输出 6 个二元组(a, apple) (b, banana) (c, cherry) (d, durian) (g, guava) (m, melon)这道题的价值在于它展示了 如何在不改变原始数据的前提下为数据附加一个用于后续分组、关联的键这是 Beam 数据流建模尤其是 Keyed 数据的第一步。三、参考实现一行 WithKeys 完成造键仓库中的参考解答 task.py 给出了完整、可运行的管道代码import apache_beam as beam with beam.Pipeline() as p: (p | beam.Create([apple, banana, cherry, durian, guava, melon]) | beam.WithKeys(lambda word: word[0:1]) | beam.LogElements())逐行拆解这段代码beam.Create([...])生成一个包含 6 个字符串元素的PCollectionbeam.WithKeys(lambda word: word[0:1])核心变换。传入一个可调用对象callableBeam 会对每个元素word调用该函数得到键word[0:1]取字符串的第一个字符用切片而非word[0]即使遇到空字符串也不会抛IndexError更稳健变换内部自动把元素包装为(键, 原值)二元组beam.LogElements()把变换结果逐条打印到标准输出方便在本地直跑或 Playground 中观察结果。注意这里使用的是有参形式传入lambda这也是本 Kata 的标准做法WithKeys同样支持传入常量键见下文第五节。运行该脚本后终端会依次打印 6 个与题目要求完全一致的元组。四、自动判题测试如何验证你的答案每个 Kata 都配有隐藏的单元测试 tests/test_task.py。该测试利用test_helper提供的test_is_not_empty和get_file_output两个工具前者检查程序输出非空后者读取task.py的实际运行输出并逐一断言以下 6 个字符串全部出现在输出中answers [(a, apple), (b, banana), (c, cherry), (d, durian), (g, guava), (m, melon)] for num in answers: self.assertIn(num, output, Incorrect output. Convert into a KV by its first letter and itself.)这套测试不仅用于 Kata 闯关的即时反馈也体现了 Beam 管道的标准验证方式给定固定输入用断言检查输出集合。在实际项目中你可以把assert_that/equal_to等 Beam 测试原语apache_beam.testing以同样思路用于管道级单元测试。五、源码原理WithKeys 在 SDK 中是如何实现的WithKeys并非一个独立的重量级PTransform类而是 SDK 中一个基于Map的轻量函数变换定义在 sdks/python/apache_beam/transforms/util.pyptransform_fn def WithKeys(pcoll, k, *args, **kwargs): PTransform that takes a PCollection, and either a constant key or a callable, and returns a PCollection of (K, V), where each of the values in the input PCollection has been paired with either the constant key or a key computed from the value. ... if callable(k): if fn_takes_side_inputs(k): ... return pcoll | Map(...) Map( lambda v, *args, **kwargs: (k(v, *args, **kwargs), v), *args, **kwargs) ... return pcoll | Map(...) Map(lambda v: (k(v), v)) return pcoll | Map(...) Map(lambda v: (k, v))从源码可以提炼出三个关键实现事实两种键模式WithKeys(pcoll, k, *args, **kwargs)的第二个参数k有两种合法形态——常量键k不是 callable 时所有元素共享同一个键实现为Map(lambda v: (k, v))计算键k是 callable 时每个元素的键由函数计算得出实现为Map(lambda v: (k(v), v))。底层就是 Map无论哪种模式最终都展开为Map变换因此WithKeys天然具备Map的并行、分布式执行特性也意味着它不会改变元素数量、不引入 shuffle是一个逐元素映射操作。支持带参 callable 与 SideInput源码通过fn_takes_side_inputs同文件 util.py检查函数签名是否接收额外参数。若k需要位置参数或关键字参数且这些参数均为AsSideInput形式WithKeys会把这些参数透传给Map从而支持基于侧输入Side Input动态计算键。官方测试覆盖的四种用法SDK 自带的 util_test.py 中WithKeysTest类用 4 个测试用例锁定了上述行为可作为学习WithKeys全部用法的活教材测试方法传入的k期望输出说明test_constant_kk常量[(k, 1), (k, 2), (k, 3)]所有元素共用一个常量键test_callable_klambda x: x * x[(1, 1), (4, 2), (9, 3)]由元素计算键值保持为原元素test_args_kwargs_k静态方法_test_args_kwargs_fn 位置/关键字参数[(1, 1), (3, 2), (5, 3)]callable 可携带额外静态参数test_sideinputslambda x, the_list, the_singleton: ...AsList/AsSingleton[(17, 1), (18, 2), (19, 3)]键的计算可依赖侧输入集合其中test_sideinputs展示了进阶用法键值可以由AsList将PCollection视为列表与AsSingleton将PCollection视为单值等侧输入动态决定适合键的规则随数据变化的场景。六、实战延伸WithKeys 的典型下游用法掌握了WithKeys之后最自然的下一步就是把生成的 KV 交给按键分组的变换这正是 Beam 中大量聚合逻辑的起点with beam.Pipeline() as p: (p | beam.Create([apple, banana, cherry, durian, guava, melon]) | beam.WithKeys(lambda word: word[0:1]) # 按首字母造键 | beam.GroupByKey() # 按键分组 | beam.Map(lambda kv: (kv[0], list(kv[1]))) # 整理分组结果 | beam.LogElements())例如本节 Kata 的数据经过GroupByKey后会得到(a, [apple])、(b, [banana])、(c, [cherry])、(d, [durian])、(g, [guava])、(m, [melon])这样的分组结果为后续CombinePerKey聚合、按用户 ID 关联事件、按类别统计等场景铺路。仓库中同样基于该主题的 Java 参考实现见 learning/beamdoc/WithKeysExample.javaJava/Kotlin 版 Kata 位于 learning/katas/java/Common Transforms、learning/katas/kotlin/Common Transforms可以对照学习各语言 API 的对应关系。需要留意的是WithKeys与ParDo的关系ParDo可以一次性完成计算键 转换值两步例如beam.ParDo中返回(k, v)而WithKeys刻意只做造键把值的转换留给后续变换从而让管道意图更清晰、更易复用。如果你的目标仅仅是把元素变成 KVWithKeys是语义最贴切的工具如果还要同时改写值则可以直接使用Map/ParDo。七、小结通过本 Kata 你可以掌握 Apache Beam Python SDK 中WithKeys变换的核心用法常量键与计算键两种形态、与Map的关系、对带参 callable 与侧输入的支持以及它与GroupByKey等分组变换的组合方式。官方实现 util.py 与测试 util_test.py 是深入理解其行为的权威参考而 task.py 与 tests/test_task.py 则是可直接运行、可直接验证的最小范例——把这两份文件结合起来阅读即可在几分钟内把WithKeys纳入你的 Beam 工具箱。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Katas 实战用 WithKeys 将 PCollection 元素转换为 KV 键值对Apache Beam Java Katas 实战用 WithKeys 将 PCollection 元素转换为 KV 键值对 导读 WithKeys 是 Ap大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战使用 WithKeys 为 PCollection 元素附加键值KVApache Beam Kotlin Kata 实战使用 WithKeys 为 PCollection 元素附加键值KV 导读 本文围绕 Apache B大数据批处理流处理数据工程Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和 导读 本文以 Apache Beam 仓库中 learni大数据批处理流处理数据工程上一篇lefthook 远程配置 ref 参数详解锁定分支与标签、源码级工作流与最佳实践下一篇scikit-learn 项目入门指南从安装依赖、运行测试到参与开发基于 README.rst 全文解读创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表