☰
Apache Beam Python Sample 聚合转换详解:FixedSizeGlobally 与 FixedSizePerKey 无放回随机抽样实战
2026/10/12 3:12:24 网站建设 项目流程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

导读

Sample是 Apache Beam Python SDK 提供的一组聚合(Aggregation)转换,用于从PCollection中随机抽取固定数量的元素,或从键值对集合中按 key 分别抽取固定数量的关联值。本文以官方文档 sample.md 为核心骨架,结合仓库源码(combiners.py、示例代码)深入讲解Sample.FixedSizeGlobally与Sample.FixedSizePerKey的用法、底层实现原理与测试验证方式,帮助你快速在批处理与流处理管道中完成随机抽样任务。


一、Sample 转换是什么

Sample位于apache_beam.transforms.combiners模块,官方文档对其定位如下:

Transforms for taking samples of the elements in a collection, or samples of the values associated with each key in a collection of key-value pairs.

即:它从集合中抽取元素样本,或从键值对集合中按 key 抽取对应值的样本。它的核心特征是无放回随机抽样(sampling n elements without replacement)——同一元素不会在结果中重复出现,且抽样结果具有随机性。

从源码看,Sample类定义在 combiners.py:

class Sample(object): """Combiners for sampling n elements without replacement.""" class FixedSizeGlobally(CombinerWithoutDefaults): """Sample n elements from the input PCollection without replacement.""" class FixedSizePerKey(ptransform.PTransform): """Sample n elements associated with each key without replacement."""

它提供两个公开转换,分别对应文档中的两个示例:

转换适用输入输出
Sample.FixedSizeGlobally(n)整个PCollection单元素PCollection,值为包含 n 个元素的List
Sample.FixedSizePerKey(n)KV<K, V>键值对PCollection每个 key 对应一个包含最多 n 个元素的List

类型注解(combiners.py)也印证了这一点:FixedSizeGlobally的输入类型为T、输出类型为List[T];FixedSizePerKey的输入类型为Tuple[K, V]、输出类型为Tuple[K, List[V]]。


二、Example 1:从整个 PCollection 随机抽样(FixedSizeGlobally)

官方文档的第一个示例演示:创建一个PCollection后,使用Sample.FixedSizeGlobally()从整个集合中获取固定大小的随机样本。

对应的完整可运行代码位于仓库 sample_fixed_size_globally.py:

import apache_beam as beam with beam.Pipeline() as pipeline: sample = ( pipeline | 'Create produce' >> beam.Create([ '🍓 Strawberry', '🥕 Carrot', '🍆 Eggplant', '🍅 Tomato', '🥔 Potato', ]) | 'Sample N elements' >> beam.combiners.Sample.FixedSizeGlobally(3) | beam.Map(print))

运行逻辑说明:

  1. beam.Create创建包含 5 个元素的PCollection;
  2. beam.combiners.Sample.FixedSizeGlobally(3)从这 5 个元素中无放回随机抽取 3 个;
  3. beam.Map(print)将结果输出到控制台。

运行结果形如(因为抽样随机,每次输出的具体元素可能不同):

['🥕 Carrot', '🍆 Eggplant', '🍅 Tomato']

注意:虽然元素内容随机,但输出列表中元素个数始终等于 n(3)。仓库中的测试 sample_test.py 正是用这个不变量做断言:

def check_sample(actual): # The sampled elements are non-deterministic, so check the sample size. assert_matches_stdout(actual, expected, lambda elements: len(elements))

底层实现:SampleCombineFn

FixedSizeGlobally的expand方法(combiners.py)内部将整个集合交给CombineGlobally(SampleCombineFn(n))聚合:

def expand(self, pcoll): if self.has_defaults: return pcoll | core.CombineGlobally(SampleCombineFn(self._n)) else: return pcoll | core.CombineGlobally( SampleCombineFn(self._n)).without_defaults()

而真正的抽样逻辑封装在SampleCombineFn(combiners.py)中,其巧妙之处在于复用TopCombineFn+ 随机数键:

class SampleCombineFn(core.CombineFn): def __init__(self, n): self._top_combiner = TopCombineFn(n) def add_input(self, heap, element): # Before passing elements to the Top combiner, we pair them with random # numbers. The elements with the n largest random number "keys" will be # selected for the output. return self._top_combiner.add_input(heap, (random.random(), element)) def extract_output(self, heap): # Here we strip off the random number keys we added in add_input. return [e for _, e in self._top_combiner.extract_output(heap)]

抽样原理可以概括为三步:

  1. 随机打标:每个元素在进入TopCombineFn前,先与一个random.random()生成的随机数配对;
  2. 取 Top-n:TopCombineFn(n)使用堆(heapq)维护随机数最大的 n 个键值对,从而等价于随机选出 n 个元素(combiners.py);
  3. 剥离随机数:extract_output时去掉随机数键,仅返回原始元素列表。

由于每个元素获得独立随机数,天然实现无放回抽样,且整体抽样概率均匀;同时借助堆的数据结构,内存占用被限制在 O(n) 级别,不会随输入规模线性增长。


三、Example 2:按 key 分别随机抽样(FixedSizePerKey)

官方文档的第二个示例演示:对KV<K, V>键值对集合使用Sample.FixedSizePerKey(),为每个唯一的 key 获取固定大小的随机样本。

对应的完整可运行代码位于仓库 sample_fixed_size_per_key.py:

import apache_beam as beam with beam.Pipeline() as pipeline: samples_per_key = ( pipeline | 'Create produce' >> beam.Create([ ('spring', '🍓'), ('spring', '🥕'), ('spring', '🍆'), ('spring', '🍅'), ('summer', '🥕'), ('summer', '🍅'), ('summer', '🌽'), ('fall', '🥕'), ('fall', '🍅'), ('winter', '🍆'), ]) | 'Samples per key' >> beam.combiners.Sample.FixedSizePerKey(3) | beam.Map(print))

运行逻辑说明:

  1. beam.Create创建包含 4 个季节 key(spring/summer/fall/winter)共 10 个键值对的PCollection;
  2. beam.combiners.Sample.FixedSizePerKey(3)对每个 key 分别执行无放回随机抽样,最多抽取 3 个值;
  3. beam.Map(print)输出形如(key, [values...])的结果。

运行结果形如(抽样随机,内容可能变化,但每个 key 的样本个数受限于输入数量):

('spring', ['🍓', '🥕', '🍆']) ('summer', ['🥕', '🍅', '🌽']) ('fall', ['🥕', '🍅']) ('winter', ['🍆'])

注意一个关键细节:n是目标样本数的上限。当某个 key 的关联值数量少于 n 时(例如上面fall只有 2 个值、winter只有 1 个值),返回的就是该 key 的全部值,不会凭空补足到 3 个。仓库测试 sample_test.py 用(key, 样本个数)校验了这一行为。

底层实现:CombinePerKey

FixedSizePerKey的expand方法(combiners.py)将键值对集合交给CombinePerKey(SampleCombineFn(n)):

def expand(self, pcoll): return pcoll | core.CombinePerKey(SampleCombineFn(self._n))

CombinePerKey(定义于 core.py)会先识别输入中具有相同 key 的值集合,再对每个 key 分别应用CombineFn进行归并——因此每个 key 的抽样彼此独立,使用与全局抽样完全相同的SampleCombineFn实现,保证了行为一致性。


四、参数说明与注意事项

参数n

两个转换都只接受一个必填参数n:

参数类型含义说明
nint目标样本数当元素总数 ≥ n 时,输出恰好 n 个;当元素总数 < n 时,输出全部元素
  • 该参数在display_data中被登记为{'n': self._n}(combiners.py),可在作业可视化面板中查看;
  • 转换的default_label为FixedSizeGlobally(n)或FixedSizePerKey(n)(combiners.py),便于在数据流图中识别。

空输入与全局聚合的默认值行为

FixedSizeGlobally继承自CombinerWithoutDefaults,其内部CombineGlobally在空输入时如何处理,取决于管道配置。从 core.py 的实现看:

  • 使用without_defaults()时,空输入产出空PCollection(无输出);
  • 使用默认模式且窗口不是全局窗口(如固定时间窗口)时,需要显式指定默认值行为,否则可能抛出ValueError,提示改用without_defaults()或as_singleton_view()。

FixedSizePerKey则天然不受此影响:每个 key 独立聚合,空输入只会得到空结果集。

抽样结果的随机性与确定性

  • 抽样结果非确定性:依赖random.random(),每次运行抽取的元素可能不同;
  • 测试与下游逻辑应基于“样本大小”而非“具体样本内容”做断言(参考 sample_test.py 的注释 "The sampled elements are non-deterministic, so check the sample size.");
  • 若需要可复现结果,可在管道层面自行管理随机种子,但SampleCombineFn本身不提供种子参数。

五、源码测试验证

仓库通过两级测试验证Sample转换的正确性:

1. 示例级测试sample_test.py

  • test_sample_fixed_size_globally:断言全局抽样结果长度恒为 3;
  • test_sample_fixed_size_per_key:断言每个 key 的样本个数不超过 3,且与输入数量匹配;
  • 使用assert_matches_stdout结合TestPipeline在真实管道中运行。

2. 单元级测试combiners_test.py

  • test_global_sample:对[1, 1, 2, 2]输入执行FixedSizeGlobally(3),断言sorted(actual[0])必为[1, 1, 2]或[1, 2, 2](即必须无放回且数量为 3);同时验证带时间戳窗口下without_defaults()路径;
  • test_per_key_sample:对 9 个 key 各 4 个值的输入执行FixedSizePerKey(3),断言每个 key 恰好输出 3 个样本,且其中 1 和 2 的数量各为 1 或 2,证明无放回且随机。

此外,combiners_test.py 还将Sample.FixedSizePerKey与Sample.FixedSizeGlobally纳入分布式(dist)场景的逐 key 测试,覆盖多 runner 下的行为一致性。


六、典型应用场景

结合Sample的语义,其典型用途包括:

  • 数据探索与采样:在建模前从海量数据中随机抽取固定比例/数量的样本,降低下游处理与可视化成本;
  • 分层抽样:对带类别 key(如地区、用户分组、季节)的键值对数据,按类别各自抽取代表性样本,保证每类都有覆盖;
  • 负载均衡/压测准备:从消息流或日志中随机抽取 n 条用于本地调试、压测或审查;
  • 与 Top 配合:官方文档在 “Related transforms” 中将Top列为关联转换——Sample用于随机抽样,而 Top 文档 用于取最大/最小元素,二者组合可完成“先抽样再取极值”的近似分析流程。

七、小结

Sample转换是 Apache Beam Python SDK 中实现随机抽样的标准工具:

  • Sample.FixedSizeGlobally(n):从整个集合无放回抽取 n 个元素;
  • Sample.FixedSizePerKey(n):按 key 分别无放回抽取至多 n 个关联值;
  • 底层由SampleCombineFn(combiners.py)基于“随机数键 + Top 堆”实现,内存高效(O(n))且抽样均匀;
  • 抽样结果非确定性,测试应基于样本数量断言。

如需继续深入,可阅读同目录下的 Top 文档、组合器基类CombineFn的实现(core.py),以及示例代码所在的 aggregation 目录。

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:Fay框架API文档暗黑模式对比度调整:符合标准
下一篇:jellyfin-ffmpeg vs 官方FFmpeg:5大独家增强功能深度对比

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询