【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文围绕 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.
也就是说,它解决两类问题:
- 从普通集合(
PCollection<T>)中抽取若干元素样本; - 从键值对集合(
PCollection<KV<K, V>>)中,为每一个不同的 Key抽取若干关联值的样本。
在 Java SDK 中,Sample是一个静态工具类(public class Sample),它对外提供了一系列静态工厂方法,分别返回不同的PTransform或CombineFn,直接apply到PCollection上即可使用。
二、核心 API 全景
从源码 Sample.java 可以看到,Sample类共暴露了 6 个静态工厂方法,按功能可分为"变换"与"CombineFn"两类:
| 方法签名 | 返回类型 | 用途 | 采样性质 |
|---|---|---|---|
Sample.any(long limit) | PTransform<PCollection<T>, PCollection<T>> | 返回输入集合中至多 limit 个元素 | 非均匀,速度快,无随机性保证 |
Sample.fixedSizeGlobally(int sampleSize) | PTransform<PCollection<T>, PCollection<Iterable<T>>> | 全局均匀随机抽取sampleSize个元素 | 均匀随机 |
Sample.fixedSizePerKey(int sampleSize) | PTransform<PCollection<KV<K,V>>, PCollection<KV<K,Iterable<V>>>> | 为每个 Key均匀随机抽取sampleSize个关联值 | 均匀随机 |
Sample.combineFn(int sampleSize) | CombineFn<T, ?, Iterable<T>> | 供Combine手动使用的 CombineFn,计算固定大小的均匀样本 | 均匀随机 |
Sample.anyCombineFn(int sampleSize) | CombineFn<T, ?, Iterable<T>> | 供Combine手动使用的 CombineFn,计算固定大小的可能非均匀样本 | 非均匀 |
Sample.anyValueCombineFn() | CombineFn<T, ?, 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 是水果 PCollection<KV<String, 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 个关联值 PCollection<KV<String, Iterable<String>>> result = pairs.apply(Sample.fixedSizePerKey(2)); result.apply(ParDo.of(new LogOutput<>( "PCollection pairs after Sample.fixedSizePerKey transform: "))); pipeline.run(); } static class LogOutput<T> extends DoFn<T, 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的输出是PCollection<KV<K, Iterable<V>>>,即每个 Key 对应一个值的 Iterable 集合,而不是打平后的单个元素;后续如果需要继续处理每个样本值,需要配合Flatten或ParDo展开。
四、三种主要变换逐个拆解
4.1Sample.any(long limit):最快的"随手抓取"
PCollection<String> input = ...; PCollection<String> output = input.apply(Sample.<String>any(100));any的行为在 Sample.java 中有明确说明:
- 输入
PCollection<T>,输出PCollection<T>(类型不变,且元素不打包成 Iterable),最多包含limit个元素; - 如果
limit大于或等于输入集合大小,则全部元素都会被选中; - 它不做均匀随机化,只保证"取到 limit 个元素",因此执行更快。
从实现看,Any变换内部是两步组合(Sample.java#L161-L165):
@Override public PCollection<T> expand(PCollection<T> in) { return in.apply(Combine.globally(new SampleAnyCombineFn<T>(limit)).withoutDefaults()) .apply(Flatten.iterables()); }即先做一次全局Combine(用withoutDefaults(),空输入不产出默认值),再通过Flatten.iterables()把样本从 Iterable 打平回单个元素。其核心SampleAnyCombineFn的实现非常直白(Sample.java#L217-L259):累积器就是一个ArrayList,addInput时只要列表长度还没到limit就追加,先到先得,不做任何随机挑选,这正是"快但非均匀"的来源。
典型适用场景:只需要"随便挑 N 条看看",例如对超大集合做快速人工抽检、先取一部分数据做局部调试。对抽样均匀性有要求时请勿使用。
4.2Sample.fixedSizeGlobally(int sampleSize):全局均匀抽样
PCollection<String> pc = ...; PCollection<Iterable<String>> sampleOfSize10 = pc.apply(Sample.fixedSizeGlobally(10));这是"真正均匀随机"的抽样方案(Sample.java#L94-L117):
- 输出类型变为
PCollection<Iterable<T>>:整个输出集合中只有一个元素,这个元素是包含sampleSize个抽样结果的Iterable; - 如果输入集合元素数少于
sampleSize,输出 Iterable 将包含全部输入元素; - 关键限制(源码注释明确给出):输出集合的所有元素必须能放进单个 worker 机器的内存,且该操作不并行执行。
适用场景:样本量远小于全量数据、且需要把整份样本集中到一处做后续分析(如一次性加载到内存、写入单文件)时使用。因为结果集中在一个元素里,对超大输入要谨慎,避免单机内存溢出。
4.3Sample.fixedSizePerKey(int sampleSize):按 Key 分组均匀抽样
PCollection<KV<String, Integer>> pc = ...; PCollection<KV<String, Iterable<Integer>>> sampleOfSize10PerKey = pc.apply(Sample.<String, Integer>fixedSizePerKey(10));针对KV输入(Sample.java#L119-L144):
- 输入
PCollection<KV<K, V>>,输出PCollection<KV<K, Iterable<V>>>; - 对每个不同的 Key,均匀随机抽取至多
sampleSize个关联值; - 若某个 Key 关联值不足
sampleSize个,则该 Key 输出其全部关联值。
实现上它直接复用了FixedSizedSampleFn,通过Combine.perKey让每个 Key 独立完成采样(Sample.java#L204-L207),天然具备并行性和扩展性,适合"按用户抽样""按品类抽样"这类分组抽检需求。
五、源码级原理:均匀抽样如何实现
fixedSizeGlobally/fixedSizePerKey之所以能保证均匀随机,关键在于它们复用了Top变换中的有界堆(Bounded Heap)数据结构。核心逻辑在Sample.FixedSizedSampleFn(Sample.java#L296-L362):
public static class FixedSizedSampleFn<T> extends CombineFn<T, Top.BoundedHeap<KV<Integer, T>, SerializableComparator<KV<Integer, T>>>, Iterable<T>> { private final int sampleSize; private final Top.TopCombineFn<KV<Integer, T>, SerializableComparator<KV<Integer, 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.BoundedHeap<KV<Integer, T>, SerializableComparator<KV<Integer, T>>> addInput( Top.BoundedHeap<KV<Integer, T>, SerializableComparator<KV<Integer, T>>> accumulator, T input) { accumulator.addInput(KV.of(rand.nextInt(), input)); return accumulator; } ... }这个算法的巧妙之处在于用随机数做排序键:
- 每来一个元素,就为它生成一个
rand.nextInt()随机整数作为"随机键",包装成KV<Integer, T>; - 交给
Top.TopCombineFn(源码见sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Top.java,即"Top 变换找最大 N 个元素"的 Combine 实现)按随机键排序,保留随机键最大的sampleSize个; - 因为随机键是均匀分布的,保留下来的元素集合在数学上就是输入的一个均匀随机子集;
extractOutput时丢弃随机键,只把元素本身放入输出List(Sample.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 theCombinetransform):
PCollection<T> input = ...; PCollection<Iterable<T>> sampled = input.apply( Combine.globally(Sample.<T>combineFn(100)));六、边界行为与参数约束
结合 Sample.java 的构造器校验与测试用例,可以总结出以下确定性行为:
| 场景 | 行为 |
|---|---|
sampleSize/limit为负数 | 抛异常:any抛IllegalArgumentException(checkArgument(limit >= 0, ...),Sample.java#L157);FixedSizedSampleFn抛IllegalArgumentException("sample size must be >= 0")(Sample.java#L304-L307) |
sampleSize/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用不同limit(0、1、半量、全量、超量)验证Sample.any的输出恰好是min(limit, 输入大小)个且是输入的子集(SampleTest.java#L71-L152);testSampleAny、testSampleAnyEmpty、testSampleAnyZero、testSampleAnyInsufficientElements、testSampleAnyNegative覆盖any的窗口、空输入、零值、不足与负数场景(SampleTest.java#L221-L294);testSample、testSampleEmpty、testSampleZero、testSampleInsufficientElements、testSampleNegative、testSampleMultiplicity覆盖fixedSizeGlobally的完整边界(SampleTest.java#L296-L368);testDisplayData验证采样大小会以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与有状态处理逻辑,为高级应用提供了扩展空间。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Java Keys 变换全面指南:从 KV 键值对集合中提取键
Apache Beam Java Keys 变换全面指南:从 KV 键值对集合中提取键 Apache Beam 的 Keys 变换属于 Java SDK 的 E
Apache Beam Java Sample 变换:从 PCollection 中随机采样的完整实战指南
Apache Beam Java Sample 变换:从 PCollection 中随机采样的完整实战指南 Apache Beam 的 Sample 变换(位于
大数据批处理流处理数据工程react-router-redux源码中的ES6特性:箭头函数与解构赋值
react router redux源码中的ES6特性:箭头函数与解构赋值 你是否在阅读React相关项目源码时,被大量简洁却陌生的语法搞得晕头转向?本文将带你
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考