【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
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))运行逻辑说明:
beam.Create创建包含 5 个元素的PCollection;beam.combiners.Sample.FixedSizeGlobally(3)从这 5 个元素中无放回随机抽取 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)]抽样原理可以概括为三步:
- 随机打标:每个元素在进入
TopCombineFn前,先与一个random.random()生成的随机数配对; - 取 Top-n:
TopCombineFn(n)使用堆(heapq)维护随机数最大的 n 个键值对,从而等价于随机选出 n 个元素(combiners.py); - 剥离随机数:
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))运行逻辑说明:
beam.Create创建包含 4 个季节 key(spring/summer/fall/winter)共 10 个键值对的PCollection;beam.combiners.Sample.FixedSizePerKey(3)对每个 key 分别执行无放回随机抽样,最多抽取 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:
| 参数 | 类型 | 含义 | 说明 |
|---|---|---|---|
n | int | 目标样本数 | 当元素总数 ≥ 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.
相关推荐
Apache Beam Java SDK Sample 变换详解:全局与按 Key 随机采样实战
Apache Beam Java SDK Sample 变换详解:全局与按 Key 随机采样实战 导读 Sample 是 Apache Beam Java SD
批处理流处理大数据Turf.js 随机抽样指南:用 @turf/sample 从 FeatureCollection 中无放回地随机选取要素
Turf.js 随机抽样指南:用 @turf/sample 从 FeatureCollection 中无放回地随机选取要素 @turf/sample 是 Tur
数据分析Apache Beam Java Sample 变换:从 PCollection 中随机采样的完整实战指南
Apache Beam Java Sample 变换:从 PCollection 中随机采样的完整实战指南 Apache Beam 的 Sample 变换(位于
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考