- 流处理
- 后端
- 大数据
【免费下载链接】storm
Apache Storm
本篇技术指南以 Apache Storm 的 Trident 高级抽象 API 为核心,系统讲解其"Stream + 批次"数据模型、分区本地操作、重分区、聚合、分组流以及合并与连接等五类核心操作。读完本文,你将掌握 Trident 中 Function、Filter、三种 Aggregator 接口的用法与差异,理解 partitionAggregate、aggregate、persistentAggregate 的执行语义,并能结合源码与示例(如 TridentWordCount)在真实拓扑中正确编排这些操作。
Trident 的核心数据模型:Stream 与批次
Trident 的核心数据模型是Stream(流),它以一系列批次(batch)为单位被处理。一条 Stream 会被分布到集群的各个节点上,对 Stream 施加的每个操作都会在每个分区(partition)上并行执行。这种"批处理 + 并行分区"的模型,让 Trident 既能获得 Storm 的实时流处理能力,又能像批处理框架一样提供恰好一次(exactly-once)语义与高吞吐的聚合计算。
在这个模型之上,Trident API 一共定义了五类操作,理解它们的分类是掌握整个 API 的关键:
| 操作类别 | 是否产生网络传输 | 典型方法 |
|---|---|---|
| 1. 分区本地操作(Partition-local operations) | 否,仅在各分区内独立执行 | each(Function/Filter)、partitionAggregate、project、stateQuery、partitionPersist |
| 2. 重分区操作(Repartitioning operations) | 是,会重新分布 tuple 但不改变内容 | shuffle、broadcast、partitionBy、global、batchGlobal、partition |
| 3. 聚合操作(Aggregation operations) | 是,操作过程中伴随网络传输 | aggregate、persistentAggregate |
| 4. 分组流操作(Operations on grouped streams) | 由groupBy引发的重分区产生 | groupBy之后的聚合、persistentAggregate |
| 5. 合并与连接(Merges and joins) | 是,涉及多个 Stream 的组合 | TridentTopology#merge、TridentTopology#join |
下文将按这五类逐一深入。
分区本地操作(Partition-local operations)
分区本地操作不涉及任何网络传输,每个批次的分区被独立处理。它们是 Trident 中最基础、最常用的一类操作。
Functions:为 tuple 追加新字段
一个Function接收一组输入字段,并发射零个或多个 tuple作为输出。输出 tuple 的字段会被追加到原始输入 tuple 上:
- 如果 Function 没有发射任何 tuple,则原始输入 tuple 被过滤掉;
- 否则,每个输出 tuple 都会复制一份原始输入 tuple。
假设有如下 Function:
public class MyFunction extends BaseFunction { public void execute(TridentTuple tuple, TridentCollector collector) { for(int i=0; i < tuple.getInteger(0); i++) { collector.emit(new Values(i)); } } }再假设变量mystream拥有字段["a", "b", "c"],其中包含以下 tuple:
[1, 2, 3] [4, 1, 6] [3, 0, 8]执行如下代码:
mystream.each(new Fields("b"), new MyFunction(), new Fields("d"));each会把字段"b"的值作为输入传给MyFunction。第一个 tuple 的b=2会发射0、1两个值,第二个 tuple 的b=1发射0,第三个 tuple 的b=0不发射任何值(因而被过滤)。最终结果字段变为["a", "b", "c", "d"]:
[1, 2, 3, 0] [1, 2, 3, 1] [4, 1, 6, 0]在源码中,Stream.each(Fields inputFields, Function function, Fields functionFields)位于 Stream.java,它会把该函数包装进一个ProcessorNode加入拓扑;另有便捷重载each(Function, Fields)与each(Fields, Filter)分别对应无输入投影与过滤场景(见 Stream.java)。
Filters:按谓词保留或剔除 tuple
Filter接收一个 tuple,决定是否保留它。假设有如下过滤器:
public class MyFilter extends BaseFilter { public boolean isKeep(TridentTuple tuple) { return tuple.getInteger(0) == 1 && tuple.getInteger(1) == 2; } }对字段为["a", "b", "c"]的以下 tuple:
[1, 2, 3] [2, 1, 1] [2, 3, 4]执行:
mystream.each(new Fields("b", "a"), new MyFilter());注意这里输入字段顺序是("b", "a"),因此getInteger(0)取到的是b、getInteger(1)取到的是a。只有第二个 tuple 满足b==1 && a==2,结果仅保留:
[2, 1, 1]从实现上看,each(Fields, Filter)会把 Filter 包装为FilterExecutor(见 Stream.java),过滤本质上就是"发射零个字段的 Function"在语义上的特化:保留则继续传递,剔除则该 tuple 被过滤掉。
partitionAggregate:在每个分区上聚合,tuple 会被替换
partitionAggregate对每个批次分区运行一个聚合函数。与 Function 不同,partitionAggregate 发射出的 tuple 会替换掉输入给它的 tuple,而不是追加。考虑这个例子:
mystream.partitionAggregate(new Fields("b"), new Sum(), new Fields("sum"));假设输入流字段为["a", "b"],分区如下:
Partition 0: ["a", 1] ["b", 2] Partition 1: ["a", 3] ["c", 8] Partition 2: ["e", 1] ["d", 9] ["d", 10]输出流将只含一个字段"sum",每个分区输出一个 tuple:
Partition 0: [3] Partition 1: [11] Partition 2: [20]Stream提供了针对三种聚合器接口的重载方法partitionAggregate(Aggregator|CombinerAggregator|ReducerAggregator, Fields)(见 Stream.java),内部统一通过chainedAgg()链式声明器组装。
三种聚合器接口:CombinerAggregator / ReducerAggregator / Aggregator
Trident 定义了三种接口用于实现聚合,它们的能力与适用场景差异很大。
CombinerAggregator(接口见 CombinerAggregator.java):
public interface CombinerAggregator<T> extends Serializable { T init(TridentTuple tuple); T combine(T val1, T val2); T zero(); }CombinerAggregator 输出单个 tuple、单个字段。它对每个输入 tuple 执行init,然后用combine把值两两合并直到只剩一个值;若分区内没有 tuple,则发射zero()的结果。例如内置的Count实现:
public class Count implements CombinerAggregator<Long> { public Long init(TridentTuple tuple) { return 1L; } public Long combine(Long val1, Long val2) { return val1 + val2; } public Long zero() { return 0L; } }CombinerAggregator 的真正优势体现在与aggregate方法配合使用时:Trident 会自动优化计算,在网络传输 tuple 之前先做部分聚合(pre-aggregation),大幅降低网络开销。
ReducerAggregator(接口见 ReducerAggregator.java):
public interface ReducerAggregator<T> extends Serializable { T init(); T reduce(T curr, TridentTuple tuple); }ReducerAggregator 先通过init()产生一个初始值,然后对每个输入 tuple 迭代调用reduce,最终输出单个 tuple、单个值。用 ReducerAggregator 实现 Count:
public class Count implements ReducerAggregator<Long> { public Long init() { return 0L; } public Long reduce(Long curr, TridentTuple tuple) { return curr + 1; } }ReducerAggregator 还可以与persistentAggregate配合使用(下文会说明)。
Aggregator(接口见 Aggregator.java)是最通用的聚合接口:
public interface Aggregator<T> extends Operation { T init(Object batchId, TridentCollector collector); void aggregate(T state, TridentTuple tuple, TridentCollector collector); void complete(T state, TridentCollector collector); }Aggregator 可以发射任意数量、任意字段的 tuple,并且可以在执行期间的任意时刻发射。它的执行流程分三步:
- init:在开始处理批次前调用,返回值代表聚合状态,会传入
aggregate和complete; - aggregate:对批次分区中的每个输入 tuple 调用,可更新状态,也可选择性发射 tuple;
- complete:当
aggregate处理完该批次分区的所有 tuple 后调用。
用 Aggregator 实现 Count:
public class CountAgg extends BaseAggregator<CountState> { static class CountState { long count = 0; } public CountState init(Object batchId, TridentCollector collector) { return new CountState(); } public void aggregate(CountState state, TridentTuple tuple, TridentCollector collector) { state.count+=1; } public void complete(CountState state, TridentCollector collector) { collector.emit(new Values(state.count)); } }聚合器链(chained aggregation)
有时需要同时执行多个聚合器,这称为链式聚合(chaining):
mystream.chainedAgg() .partitionAggregate(new Count(), new Fields("count")) .partitionAggregate(new Fields("b"), new Sum(), new Fields("sum")) .chainEnd();这段代码会在每个分区上同时运行Count与Sum两个聚合器,输出单个 tuple,字段为["count", "sum"]。链式聚合的声明入口是ChainedAggregatorDeclarer(见 ChainedAggregatorDeclarer.java),groupBy后返回的GroupedStream也提供chainedAgg()(见 GroupedStream.java)。
stateQuery 与 partitionPersist:查询与更新状态源
stateQuery和partitionPersist分别用于查询与更新状态源(source of state),它们是与持久化存储交互的入口。详细用法可参见 Trident 状态文档。二者是后续实现"窗口连接"等高级模式的基石。
projection:投影保留指定字段
projection(方法名为project)只保留操作中指定的字段。若 Stream 字段为["a", "b", "c", "d"],执行:
mystream.project(new Fields("b", "d"));输出流将只包含字段["b", "d"]。源码中project通过ProjectedProcessor处理器实现(见 Stream.java)。
重分区操作(Repartitioning operations)
重分区操作会改变 tuple 在任务间的分布方式,分区的数量也可能随之改变(例如重分区后的并行度提示更高时),因此必然涉及网络传输。Trident 提供六种重分区函数:
- shuffle:使用随机轮询(round robin)算法,将 tuple 均匀地重新分布到所有目标分区;
- broadcast:每个 tuple 被复制到所有目标分区。这在 DRPC 场景下很有用——例如需要对每一份数据都做一次
stateQuery时; - partitionBy:接收一组字段,基于这组字段做语义分区——字段被哈希后对目标分区数取模,从而选定目标分区。
partitionBy保证相同的字段集合总是落到相同的目标分区; - global:所有 tuple 都被发送到同一个分区,且整个 Stream 的所有批次都选择该分区;
- batchGlobal:批次内的所有 tuple 被发送到同一个分区,但 Stream 中的不同批次可能落到不同分区;
- partition:接收一个实现了
backtype.storm.grouping.CustomStreamGrouping的自定义分区函数。
这些方法在 Stream.java 中均有直接实现,可以看到它们的底层分组语义:shuffle映射到Grouping.shuffle,broadcast映射到Grouping.all,global使用内部GlobalGrouping(以保证批次发射到固定分区),batchGlobal则使用IndexHashGrouping(0)——以批次的第一个字段(batch id)作为哈希键,这正是"同批次同分区、不同批次可不同分区"的机制来源。
聚合操作(Aggregation operations)
Trident 提供aggregate和persistentAggregate两个方法对流做聚合。区别在于:
- aggregate:对 Stream 的每个批次****独立执行聚合;
- persistentAggregate:对 Stream所有批次的所有 tuple做聚合,并把结果存储到状态源中。
aggregate在 Stream 上执行的是全局聚合,其执行策略取决于聚合器类型:
- 使用ReducerAggregator 或 Aggregator时:Stream 会先被重分区为单个分区,然后在该分区上运行聚合函数;
- 使用CombinerAggregator时:Trident 会先计算每个分区的部分聚合,再重分区到单个分区,最后在网络传输之后完成聚合。
因此CombinerAggregator 效率远高于另外两种,应尽可能优先使用。例如对批次求全局计数:
mystream.aggregate(new Count(), new Fields("count"));与partitionAggregate一样,aggregate的聚合器也可以链式组合。但要注意:如果把 CombinerAggregator 与非 CombinerAggregator 链在一起,Trident 将无法进行部分聚合优化。
persistentAggregate的用法详见 Trident 状态文档。从源码可以观察到它的实现策略差异(见 GroupedStream.java):
- 对CombinerAggregator:先做
aggregate(含部分聚合优化),再通过partitionPersist配合MapCombinerAggStateUpdater写入状态; - 对ReducerAggregator:先
partitionBy(_groupFields)按分组字段重分区,再partitionPersist配合MapReducerAggStateUpdater写入状态。
这解释了为何persistentAggregate的输出是一个TridentState对象,而不是普通 Stream。
分组流操作(Operations on grouped streams)
groupBy操作先对指定字段做partitionBy重分区,然后在每个分区内把分组字段相等的 tuple 归组。其分区与分组语义如下图所示:
在分组流上运行聚合器时,聚合将在每个组内进行,而不是针对整个批次。persistentAggregate也可以运行在GroupedStream上,此时结果会以分组字段作为 key存储在一个 MapState 中(详见 Trident 状态文档)。
与普通流一样,分组流上的聚合器也可以链式组合(GroupedStream.chainedAgg(),见 GroupedStream.java)。
实战示例:TridentWordCount
仓库中的 TridentWordCount.java 是上述 API 的综合演示,完整串联了 Function、groupBy、persistentAggregate 与 stateQuery:
TridentState wordCounts = topology.newStream("spout1", spout).parallelismHint(16).each(new Fields("sentence"), new Split(), new Fields("word")).groupBy(new Fields("word")).persistentAggregate(new MemoryMapState.Factory(), new Count(), new Fields("count")).parallelismHint(16); topology.newDRPCStream("words", drpc).each(new Fields("args"), new Split(), new Fields("word")).groupBy(new Fields( "word")).stateQuery(wordCounts, new Fields("word"), new MapGet(), new Fields("count")).each(new Fields("count"), new FilterNull()).aggregate(new Fields("count"), new Sum(), new Fields("sum"));该示例中,Split是自定义的BaseFunction(将句子按空格切词并逐词发射),Count、Sum、FilterNull、MapGet均为storm.trident.operation.builtin包内置的现成操作。第一条流水线把词频写入MemoryMapState;第二条 DRPC 流水线则查询词频并对命中结果求和——这正是"分组 + 状态持久化 + 状态查询"的典型组合。
合并与连接(Merges and joins)
API 的最后一部分是把不同的 Stream 组合在一起。
**合并(merge)**是最简单的组合方式,通过TridentTopology#merge完成:
topology.merge(stream1, stream2, stream3);Trident 会将新合并流的输出字段命名为第一个流的输出字段。
连接(join)则是另一种组合方式。标准的 SQL 式 join 需要有限输入,因此不适用于无限流;Trident 中的 join 只作用于 spout 产生的每个小批次内部。
看一个例子:流stream1包含字段["key", "val1", "val2"],流stream2包含字段["x", "val1"]:
topology.join(stream1, new Fields("key"), stream2, new Fields("x"), new Fields("key", "a", "b", "c"));这里分别用"key"和"x"作为两个流的连接字段。由于输入流的字段名可能重叠,Trident 要求显式指定新流的所有输出字段名。join 发射出的 tuple 内容按以下顺序构成:
- 先是连接字段列表。本例中,
"key"对应stream1的"key"、stream2的"x"; - 然后是所有流的所有非连接字段,按流传入
join方法的顺序排列。本例中,"a"、"b"对应stream1的"val1"、"val2","c"对应stream2的"val1"。
当 join 发生在来自不同 spout的流之间时,这些 spout 会被同步批次发射——即一个处理批次会包含来自每个 spout 的 tuple。
如何实现"窗口连接"(windowed join)
你可能想知道:如何实现类似"窗口连接"的语义——比如把 join 一侧的 tuple 与另一侧最近一个小时的 tuple 做连接?
做法是组合使用partitionPersist与stateQuery:把 join 一侧最近一小时的 tuple 按连接字段为 key,存储并轮转(rotate)在一个状态源中;然后用stateQuery按连接字段做查找,从而完成"连接"。这正是 Trident 状态文档 中状态机制的核心应用场景。
延伸阅读与源码参考
- Trident 教程:从零构建 Trident 拓扑的入门实践;
- Trident 状态文档:
stateQuery、partitionPersist、persistentAggregate与状态源(如MapState、MemoryMapState)的深入讲解; - Trident spouts:批次发射源头(如
FixedBatchSpout)与批次协调机制; - 核心 API 源码:Stream.java、GroupedStream.java、ChainedAggregatorDeclarer.java;
- 聚合器接口:CombinerAggregator.java、ReducerAggregator.java、Aggregator.java;
- 完整可运行示例:TridentWordCount.java 与 TridentReach.java,均位于 storm-starter 示例模块(
examples/storm-starter/pom.xml)。
掌握上述五类操作与三种聚合器接口,你就拥有了用 Trident 编排高吞吐、有状态、恰好一次语义实时计算流程的完整工具箱。在实际工程中,建议优先选用 CombinerAggregator 以获取部分聚合优化,并根据数据分布选择合适的重分区策略(partitionBy保证语义分组、global/batchGlobal用于全局归并),从而在正确性与吞吐之间取得最佳平衡。
- 流处理
- 后端
- 大数据
【免费下载链接】storm
Apache Storm
相关推荐
Apache Storm Trident API 完全指南:五类流操作、聚合、窗口与状态
Apache Storm Trident API 完全指南:五类流操作、聚合、窗口与状态 本文以 Apache Storm 的官方文档 docs/Trident
大数据流处理后端Apache Storm Trident API 全解析:从批处理模型到五大类流操作实战指南
Apache Storm Trident API 全解析:从批处理模型到五大类流操作实战指南 导读 :Trident 是 Apache Storm 之上的高层流
流处理后端大数据Apache Storm Trident API 全面解析:Stream 批处理操作、聚合、窗口与状态查询实战指南
Apache Storm Trident API 全面解析:Stream 批处理操作、聚合、窗口与状态查询实战指南 Trident 是 Apache Storm
后端大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考