Apache Beam Python Sample 聚合转换详解:FixedSizeGlobally 与 FixedSizePerKey 无放回随机抽样实战
发布时间:2026/10/12 3:12:45 锦皓数字建站

【免费下载链接】beamApache 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.pyclass 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 个元素的ListSample.FixedSizePerKey(n)KVK, 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.pyimport 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 个元素的PCollectionbeam.combiners.Sample.FixedSizeGlobally(3)从这 5 个元素中无放回随机抽取 3 个beam.Map(print)将结果输出到控制台。运行结果形如因为抽样随机每次输出的具体元素可能不同[ Carrot, Eggplant, Tomato]注意虽然元素内容随机但输出列表中元素个数始终等于 n3。仓库中的测试 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))底层实现SampleCombineFnFixedSizeGlobally的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()而真正的抽样逻辑封装在SampleCombineFncombiners.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-nTopCombineFn(n)使用堆heapq维护随机数最大的 n 个键值对从而等价于随机选出 n 个元素combiners.py剥离随机数extract_output时去掉随机数键仅返回原始元素列表。由于每个元素获得独立随机数天然实现无放回抽样且整体抽样概率均匀同时借助堆的数据结构内存占用被限制在 O(n) 级别不会随输入规模线性增长。三、Example 2按 key 分别随机抽样FixedSizePerKey官方文档的第二个示例演示对KVK, V键值对集合使用Sample.FixedSizePerKey()为每个唯一的 key 获取固定大小的随机样本。对应的完整可运行代码位于仓库 sample_fixed_size_per_key.pyimport 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 个季节 keyspring/summer/fall/winter共 10 个键值对的PCollectionbeam.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, 样本个数)校验了这一行为。底层实现CombinePerKeyFixedSizePerKey的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.pytest_sample_fixed_size_globally断言全局抽样结果长度恒为 3test_sample_fixed_size_per_key断言每个 key 的样本个数不超过 3且与输入数量匹配使用assert_matches_stdout结合TestPipeline在真实管道中运行。2. 单元级测试combiners_test.pytest_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 个关联值底层由SampleCombineFncombiners.py基于“随机数键 Top 堆”实现内存高效O(n)且抽样均匀抽样结果非确定性测试应基于样本数量断言。如需继续深入可阅读同目录下的 Top 文档、组合器基类CombineFn的实现core.py以及示例代码所在的 aggregation 目录。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐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 变换位于大数据批处理流处理数据工程上一篇Fay框架API文档暗黑模式对比度调整符合标准下一篇jellyfin-ffmpeg vs 官方FFmpeg5大独家增强功能深度对比创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。