
【免费下载链接】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变换展开讲解如何从PCollection中均匀随机抽取固定数量的元素、对KV键值对按 Key 分组采样以及如何用更快的非均匀采样Sample.any快速筛选数据。读完本文你将掌握Sample.fixedSizeGlobally、Sample.fixedSizePerKey、Sample.any三种核心 API 的用法、各自的输出类型与约束并能从源码层面理解其基于Combine与有界堆Bounded Heap的实现原理为数据探索、抽样质检、随机下采样等场景选择正确的采样方案。一、Sample 变换是什么在 Apache Beam 中Sample位于sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sample.java其官方文档website/www/site/content/en/documentation/transforms/java/aggregation/sample.md给出的定位非常明确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.也就是说它解决两类问题从普通集合PCollectionT中抽取若干元素样本从键值对集合PCollectionKVK, V中为每一个不同的 Key抽取若干关联值的样本。在 Java SDK 中Sample是一个静态工具类public class Sample它对外提供了一系列静态工厂方法分别返回不同的PTransform或CombineFn直接apply到PCollection上即可使用。二、核心 API 全景从源码 Sample.java 可以看到Sample类共暴露了 6 个静态工厂方法按功能可分为变换与CombineFn两类方法签名返回类型用途采样性质Sample.any(long limit)PTransformPCollectionT, PCollectionT返回输入集合中至多 limit 个元素非均匀速度快无随机性保证Sample.fixedSizeGlobally(int sampleSize)PTransformPCollectionT, PCollectionIterableT全局均匀随机抽取sampleSize个元素均匀随机Sample.fixedSizePerKey(int sampleSize)PTransformPCollectionKVK,V, PCollectionKVK,IterableV为每个 Key均匀随机抽取sampleSize个关联值均匀随机Sample.combineFn(int sampleSize)CombineFnT, ?, IterableT供Combine手动使用的 CombineFn计算固定大小的均匀样本均匀随机Sample.anyCombineFn(int sampleSize)CombineFnT, ?, IterableT供Combine手动使用的 CombineFn计算固定大小的可能非均匀样本非均匀Sample.anyValueCombineFn()CombineFnT, ?, T供Combine手动使用的 CombineFn计算单个可能非均匀的样本值非均匀源码注释对这两类采样有明确的语义区分fixedSizeGlobally(int)andfixedSizePerKey(int)compute uniformly random samples.any(long)is faster, but provides no uniformity guarantees.即带fixedSize前缀的 API 保证均匀随机性any系列只保证快而不保证均匀。这是选择 API 时最重要的判断依据。三、实战示例按 Key 抽取固定大小样本Sample文档页中的在线 Playground 示例对应仓库中的examples/java/src/main/java/org/apache/beam/examples/SampleExample.java该文件头部注释中name: Sample、tags: transforms/pairs/group表明它正是SDK_JAVA_Sample片段的来源。完整可运行代码如下package org.apache.beam.examples; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.transforms.Sample; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PCollection; public class SampleExample { public static void main(String[] args) { PipelineOptions options PipelineOptionsFactory.create(); Pipeline pipeline Pipeline.create(options); // 创建键值对集合key 是季节value 是水果 PCollectionKVString, String pairs pipeline.apply( Create.of( KV.of(fall, apple), KV.of(spring, strawberry), KV.of(winter, orange), KV.of(summer, peach), KV.of(spring, cherry), KV.of(fall, pear))); // 为每个唯一的 Key 均匀随机抽取 2 个关联值 PCollectionKVString, IterableString result pairs.apply(Sample.fixedSizePerKey(2)); result.apply(ParDo.of(new LogOutput( PCollection pairs after Sample.fixedSizePerKey transform: ))); pipeline.run(); } static class LogOutputT extends DoFnT, T { private static final Logger LOG LoggerFactory.getLogger(LogOutput.class); private final String prefix; public LogOutput(String prefix) { this.prefix prefix; } ProcessElement public void processElement(ProcessContext c) throws Exception { LOG.info(prefix c.element()); c.output(c.element()); } } }运行结果示意每个 Key 抽取 2 个值被抽中的具体值是随机的PCollection pairs after Sample.fixedSizePerKey transform: KV{fall, [apple, pear]} PCollection pairs after Sample.fixedSizePerKey transform: KV{spring, [strawberry, cherry]} PCollection pairs after Sample.fixedSizePerKey transform: KV{winter, [orange]} PCollection pairs after Sample.fixedSizePerKey transform: KV{summer, [peach]}需要注意fixedSizePerKey的输出是PCollectionKVK, IterableV即每个 Key 对应一个值的 Iterable 集合而不是打平后的单个元素后续如果需要继续处理每个样本值需要配合Flatten或ParDo展开。四、三种主要变换逐个拆解4.1Sample.any(long limit)最快的随手抓取PCollectionString input ...; PCollectionString output input.apply(Sample.Stringany(100));any的行为在 Sample.java 中有明确说明输入PCollectionT输出PCollectionT类型不变且元素不打包成 Iterable最多包含limit个元素如果limit大于或等于输入集合大小则全部元素都会被选中它不做均匀随机化只保证取到 limit 个元素因此执行更快。从实现看Any变换内部是两步组合Sample.java#L161-L165Override public PCollectionT expand(PCollectionT in) { return in.apply(Combine.globally(new SampleAnyCombineFnT(limit)).withoutDefaults()) .apply(Flatten.iterables()); }即先做一次全局Combine用withoutDefaults()空输入不产出默认值再通过Flatten.iterables()把样本从 Iterable 打平回单个元素。其核心SampleAnyCombineFn的实现非常直白Sample.java#L217-L259累积器就是一个ArrayListaddInput时只要列表长度还没到limit就追加先到先得不做任何随机挑选这正是快但非均匀的来源。典型适用场景只需要随便挑 N 条看看例如对超大集合做快速人工抽检、先取一部分数据做局部调试。对抽样均匀性有要求时请勿使用。4.2Sample.fixedSizeGlobally(int sampleSize)全局均匀抽样PCollectionString pc ...; PCollectionIterableString sampleOfSize10 pc.apply(Sample.fixedSizeGlobally(10));这是真正均匀随机的抽样方案Sample.java#L94-L117输出类型变为PCollectionIterableT整个输出集合中只有一个元素这个元素是包含sampleSize个抽样结果的Iterable如果输入集合元素数少于sampleSize输出 Iterable 将包含全部输入元素关键限制源码注释明确给出输出集合的所有元素必须能放进单个 worker 机器的内存且该操作不并行执行。适用场景样本量远小于全量数据、且需要把整份样本集中到一处做后续分析如一次性加载到内存、写入单文件时使用。因为结果集中在一个元素里对超大输入要谨慎避免单机内存溢出。4.3Sample.fixedSizePerKey(int sampleSize)按 Key 分组均匀抽样PCollectionKVString, Integer pc ...; PCollectionKVString, IterableInteger sampleOfSize10PerKey pc.apply(Sample.String, IntegerfixedSizePerKey(10));针对KV输入Sample.java#L119-L144输入PCollectionKVK, V输出PCollectionKVK, IterableV对每个不同的 Key均匀随机抽取至多sampleSize个关联值若某个 Key 关联值不足sampleSize个则该 Key 输出其全部关联值。实现上它直接复用了FixedSizedSampleFn通过Combine.perKey让每个 Key 独立完成采样Sample.java#L204-L207天然具备并行性和扩展性适合按用户抽样按品类抽样这类分组抽检需求。五、源码级原理均匀抽样如何实现fixedSizeGlobally/fixedSizePerKey之所以能保证均匀随机关键在于它们复用了Top变换中的有界堆Bounded Heap数据结构。核心逻辑在Sample.FixedSizedSampleFnSample.java#L296-L362public static class FixedSizedSampleFnT extends CombineFnT, Top.BoundedHeapKVInteger, T, SerializableComparatorKVInteger, T, IterableT { private final int sampleSize; private final Top.TopCombineFnKVInteger, T, SerializableComparatorKVInteger, T topCombineFn; private final Random rand new Random(); private FixedSizedSampleFn(int sampleSize) { if (sampleSize 0) { throw new IllegalArgumentException(sample size must be 0); } this.sampleSize sampleSize; topCombineFn new Top.TopCombineFn(sampleSize, new KV.OrderByKey()); } Override public Top.BoundedHeapKVInteger, T, SerializableComparatorKVInteger, T addInput( Top.BoundedHeapKVInteger, T, SerializableComparatorKVInteger, T accumulator, T input) { accumulator.addInput(KV.of(rand.nextInt(), input)); return accumulator; } ... }这个算法的巧妙之处在于用随机数做排序键每来一个元素就为它生成一个rand.nextInt()随机整数作为随机键包装成KVInteger, T交给Top.TopCombineFn源码见sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Top.java即Top 变换找最大 N 个元素的 Combine 实现按随机键排序保留随机键最大的sampleSize个因为随机键是均匀分布的保留下来的元素集合在数学上就是输入的一个均匀随机子集extractOutput时丢弃随机键只把元素本身放入输出ListSample.java#L335-L343。这就是均匀随机 分布式 Combine 归并可并行的来源每个分片先用有界堆选出局部 Top-N再由Combine逐层归并mergeAccumulators直接复用topCombineFn.mergeAccumulators最终合并结果依然是全局随机 Top-N。整套机制是Top变换找最大/最小 N 个元素在随机键维度上的复用两者都位于org.apache.beam.sdk.transforms包下。另外Sample还提供了可手动配合Combine与有状态处理使用的 CombineFn源码注释明确说明combineFncan also be used manually, in combination with state and with theCombinetransformPCollectionT input ...; PCollectionIterableT sampled input.apply( Combine.globally(Sample.TcombineFn(100)));六、边界行为与参数约束结合 Sample.java 的构造器校验与测试用例可以总结出以下确定性行为场景行为sampleSize/limit为负数抛异常any抛IllegalArgumentExceptioncheckArgument(limit 0, ...)Sample.java#L157FixedSizedSampleFn抛IllegalArgumentException(sample size must be 0)Sample.java#L304-L307sampleSize/limit为 0输出为空空 Iterable / 空集合不报错输入集合为空输出为空不报错Combine.globally(...).withoutDefaults()保证空输入不产出空默认值输入元素数小于 sampleSize输出包含全部输入元素输入元素数大于等于 sampleSize输出恰好sampleSize个元素fixedSize*为均匀随机any为任意挑选元素重复multiplicity重复元素可被多次抽中抽样按出现位置进行不保证去重这些边界行为在sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/SampleTest.java中都有对应的参数化测试与独立用例覆盖例如PickAnyTest用不同limit0、1、半量、全量、超量验证Sample.any的输出恰好是min(limit, 输入大小)个且是输入的子集SampleTest.java#L71-L152testSampleAny、testSampleAnyEmpty、testSampleAnyZero、testSampleAnyInsufficientElements、testSampleAnyNegative覆盖any的窗口、空输入、零值、不足与负数场景SampleTest.java#L221-L294testSample、testSampleEmpty、testSampleZero、testSampleInsufficientElements、testSampleNegative、testSampleMultiplicity覆盖fixedSizeGlobally的完整边界SampleTest.java#L296-L368testDisplayData验证采样大小会以sampleSize为键写入 DisplayData便于在监控/UI 中查看配置SampleTest.java#L375-L384。七、与其他聚合变换的关系Sample属于 Java SDK 的Aggregation聚合变换家族文档页将其与以下变换列为相关变换链接已转换为仓库根目录相对路径Top求最大/最小的 N 个元素与Sample共享底层Top.BoundedHeap实现。区别在于Top按元素本身的比较器取前 N 个Sample.fixedSize*则按随机键取前 N 个从而把找 Top-N变成了均匀采样Latest计算集合中最新的元素同样提供globally()与perKey()语义但选择标准是元素的隐式时间戳而非随机性适合取每 Key 最新一条的需求。选型速查需要均匀随机样本统计分析、抽样调查→fixedSizeGlobally/fixedSizePerKey只需要快速随便取 N 条调试、快速抽检→any需要把样本集中到单机内存进一步处理 →fixedSizeGlobally注意内存限制需要按 Key 分组抽检按用户、按品类→fixedSizePerKey需要在Combine或状态处理中内嵌采样逻辑 →combineFn/anyCombineFn/anyValueCombineFn。八、总结Apache Beam Java 的Sample变换是数据管道中抽样这一基础操作的标准实现fixedSizeGlobally与fixedSizePerKey通过随机键 Top 有界堆 Combine 归并在分布式环境下实现均匀随机抽样any则以放弃均匀性为代价换取更快的执行速度。理解其输出类型Iterable 打包 vs 打平元素、内存约束与边界行为能帮助你在数据探索、抽检和质量控制等场景中正确选用 API。更进一步其 CombineFn 版本可被灵活嵌入自定义Combine与有状态处理逻辑为高级应用提供了扩展空间。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Keys 变换全面指南从 KV 键值对集合中提取键Apache Beam Java Keys 变换全面指南从 KV 键值对集合中提取键 Apache Beam 的 Keys 变换属于 Java SDK 的 EApache Beam Java Sample 变换从 PCollection 中随机采样的完整实战指南Apache Beam Java Sample 变换从 PCollection 中随机采样的完整实战指南 Apache Beam 的 Sample 变换位于大数据批处理流处理数据工程react-router-redux源码中的ES6特性箭头函数与解构赋值react router redux源码中的ES6特性箭头函数与解构赋值 你是否在阅读React相关项目源码时被大量简洁却陌生的语法搞得晕头转向本文将带你上一篇Apache Pulsar TLS 客户端认证实战指南从证书签发到 Broker/Proxy/多语言客户端配置下一篇如何高效处理RPG Maker加密资源纯前端解密方案深度解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考