【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
Values是 Apache Beam Python SDK 中一个轻量的 Elementwise 变换(位于apache_beam.transforms.util模块),作用是从键值对(key-value pair)构成的PCollection中提取每个元素的"值"部分并丢弃"键"部分。本文以官方文档 values.md 为骨架,结合 util.py 源码实现、snippets 示例 与官方测试用例,带你掌握Values的用法、实现原理及其在真实 Beam 作业中的典型应用场景(例如从聚合结果中剥离键、计算全局均值等)。
什么是Values变换
Values接收一个由键值对组成的PCollection,对其中每个元素执行"取二元组第二个分量"的操作,输出一个仅包含值的新PCollection,键信息在此过程中被完全丢弃。
其核心行为可用一句话概括(引自官方文档):
Takes a collection of key-value pairs, and returns the value of each element.
该变换属于apache_beam.transforms.util模块下的"键值对工具"家族,与它并列的还有Keys(提取键)、KvSwap(键值互换)等变换。它是 Beam 官方文档中 Elementwise(逐元素)变换类别的一员,常用于数据清洗、聚合后处理等场景。
快速上手:完整可运行示例
官方在 values.py 中提供了一个可直接运行的"菜园植物"示例:先用beam.Create创建一组"表情符号 → 植物名称"的键值对,再通过beam.Values()提取出所有植物名称:
import apache_beam as beam with beam.Pipeline() as pipeline: plants = ( pipeline | 'Garden plants' >> beam.Create([ ('🍓', 'Strawberry'), ('🥕', 'Carrot'), ('🍆', 'Eggplant'), ('🍅', 'Tomato'), ('🥔', 'Potato'), ]) | 'Values' >> beam.Values() | beam.Map(print))运行该代码段后,控制台将输出:
Strawberry Carrot Eggplant Tomato Potato注意:Beam 的分布式执行模型中,PCollection内元素的无序性由 Runner 决定(Direct Runner 下通常保持输入顺序),因此实际输出的元素顺序可能不同,但内容集合一定是这五个植物名称。
源码剖析:Values是如何实现的
Values的实现非常简洁,位于 util.py:
@ptransform_fn @typehints.with_input_types(Tuple[K, V]) @typehints.with_output_types(V) def Values(pcoll, label='Values'): # pylint: disable=invalid-name """Produces a PCollection of second elements of 2-tuples in a PCollection.""" return pcoll | label >> MapTuple(lambda _, v: v)从中可以提炼出三个关键实现细节:
1. 函数式 PTransform 定义(@ptransform_fn)
Values没有采用传统class Xxx(PTransform)+expand()的写法,而是通过 ptransform.py 中的ptransform_fn装饰器以普通函数形式定义。装饰器会将其包装为_PTransformFnPTransform,使函数既能像pcoll.apply(beam.Values())一样被调用,也能通过管道符pcoll | beam.Values()组合,两种方式等价。
2. 类型提示(Type Hints)
with_input_types(Tuple[K, V]):声明输入必须是二元组,K、V为任意类型参数;with_output_types(V):声明输出为值的类型。
借助类型提示,Beam 可以在图构建阶段做静态类型检查。例如将Values应用于一个非二元组的PCollection时,会触发类型检查错误。与之类似的检查可见 ptransform_test.py 中GroupByKey的类型违规测试(expected Tuple[TypeVariable[K], TypeVariable[V]])。
3. 核心逻辑:MapTuple(lambda _, v: v)
MapTuple将输入二元组解包后传入 lambda,lambda _, v: v丢弃第一个分量(键,用_占位)并返回第二个分量(值)。因此Values本质上是Map的一个特化形式——等价于beam.Map(lambda kv: kv[1]),但语义更清晰、更贴合键值对处理的表达习惯。
参数说明:自定义 transform label
Values只接受一个可选参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
pcoll | PCollection | — | 输入键值对集合,第一个位置参数(管道符写法下自动注入) |
label | str | 'Values' | 该变换在图中的名称,用于日志、监控指标与执行计划展示 |
当同一 Pipeline 中多次使用Values时,建议为每次使用指定不同的 label 以提升可观测性,例如:
vals = pcoll.apply(beam.Values('vals'))这种"重命名实例"的用法在官方测试 ptransform_test.py 中有所体现(pcoll.apply(beam.Values('vals')))。
典型应用场景:聚合结果的后处理
Values最常见的实战场景是与GroupByKey、CombinePerKey等键值类变换配合,剥离键后继续对值做全局运算。官方真实示例 game_stats.py(游戏用户得分统计)展示了这一模式:先用CombinePerKey(sum)得到(user, total_score),再用Values()取出所有总分并计算全局均值:
# Get the sum of scores for each user. sum_scores = (user_scores | 'SumUsersScores' >> beam.CombinePerKey(sum)) # Extract the score from each element, and use it to find the global mean. global_mean_score = ( sum_scores | beam.Values() | beam.CombineGlobally(beam.combiners.MeanCombineFn())\ .as_singleton_view())类似地,官方测试 combiners_test.py 中,test_MeanCombineFn_combine先用Values()去掉键再计算全局均值,与CombinePerKey的"按键均值"形成对照,验证了两者的语义差异。
这种"先按键聚合、后Values剥离键做全局运算"的管道模式,在需要同时输出按 key 明细和全局统计的报表类作业中非常普遍。
相关变换:Keys与KvSwap
官方文档明确指出Values的两个相关变换(同样位于 util.py):
| 变换 | 输入 → 输出 | 源码实现 |
|---|---|---|
Keys | (K, V)→K,提取每个元素的键 | MapTuple(lambda k, _: k) |
KvSwap | (K, V)→(V, K),交换键和值 | MapTuple(lambda k, v: (v, k)) |
Values | (K, V)→V,提取每个元素的值 | MapTuple(lambda _, v: v) |
三者一一对应、彼此互补:Keys取第一分量,Values取第二分量,KvSwap则把两个分量对调。官方测试 ptransform_test.py 中test_keys_and_values与test_kv_swap分别验证了三者的行为:
def test_keys_and_values(self): with TestPipeline() as pipeline: pcoll = pipeline | 'Start' >> beam.Create([(3, 1), (2, 1), (1, 1), (3, 2), (2, 2), (3, 3)]) keys = pcoll.apply(beam.Keys('keys')) vals = pcoll.apply(beam.Values('vals')) assert_that(keys, equal_to([1, 2, 2, 3, 3, 3]), label='assert:keys') assert_that(vals, equal_to([1, 1, 1, 2, 2, 3]), label='assert:vals')使用注意事项
- 输入必须是二元组:
Values期望每个元素是(key, value)形式的二元组。若输入元素不是二元组,图构建阶段的类型检查或运行期MapTuple解包会报错。 - 键被完全丢弃:与
KvSwap(保留双方)不同,Values不保留键的任何信息。如果后续仍需要键值对应关系,应先完成键上的所有计算,再使用Values。 - 不改变元素个数:
Values是逐元素(elementwise)变换,输入几个元素就输出几个元素,不涉及聚合或分组;去重需求应交给Distinct等变换。 - 惰性执行:
Values只在 Pipeline 运行时真正执行,代码中书写管道只是描述执行计划,这与其他所有 Beam 变换一致。
小结
Values虽小,却是 Beam Python SDK 键值对处理链路中的基础构件:它通过MapTuple以极简代码实现了"丢弃键、保留值",配合类型提示保障输入合法性,并在真实项目中承担着"聚合结果剥离键做全局统计"的关键职责。结合本文给出的示例、源码定位与测试证据,你可以在自己的 Beam 管道中放心使用它,并与Keys、KvSwap组合出清晰的键值对处理流程。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Python Keys 变换详解:从键值对 PCollection 中提取 Key
Apache Beam Python Keys 变换详解:从键值对 PCollection 中提取 Key 本文围绕 Apache Beam Python SD
Apache Beam Java SDK Values 转换详解:从 KV 集合中提取值的逐元素变换
Apache Beam Java SDK Values 转换详解:从 KV 集合中提取值的逐元素变换 Values 是 Apache Beam Java SDK
大数据批处理流处理数据工程Instant 2025 年 11 月更新解读:MCP 工具升级、Explorer 数据管理与 Firebase 登录接入
Instant 2025 年 11 月更新解读:MCP 工具升级、Explorer 数据管理与 Firebase 登录接入 本篇文章以 Instant 项目 2
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考