☰
Apache Storm Trident API 完全指南:Stream 数据模型与五类核心流操作详解
2026/10/9 7:38:02 网站建设 项目流程
  • 流处理
  • 后端
  • 大数据

【免费下载链接】storm

Apache Storm

项目地址:https://gitcode.com/gh_mirrors/storm26/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,并且可以在执行期间的任意时刻发射。它的执行流程分三步:

  1. init:在开始处理批次前调用,返回值代表聚合状态,会传入aggregate和complete;
  2. aggregate:对批次分区中的每个输入 tuple 调用,可更新状态,也可选择性发射 tuple;
  3. 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 提供六种重分区函数:

  1. shuffle:使用随机轮询(round robin)算法,将 tuple 均匀地重新分布到所有目标分区;
  2. broadcast:每个 tuple 被复制到所有目标分区。这在 DRPC 场景下很有用——例如需要对每一份数据都做一次stateQuery时;
  3. partitionBy:接收一组字段,基于这组字段做语义分区——字段被哈希后对目标分区数取模,从而选定目标分区。partitionBy保证相同的字段集合总是落到相同的目标分区;
  4. global:所有 tuple 都被发送到同一个分区,且整个 Stream 的所有批次都选择该分区;
  5. batchGlobal:批次内的所有 tuple 被发送到同一个分区,但 Stream 中的不同批次可能落到不同分区;
  6. 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 内容按以下顺序构成:

  1. 先是连接字段列表。本例中,"key"对应stream1的"key"、stream2的"x";
  2. 然后是所有流的所有非连接字段,按流传入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

项目地址:https://gitcode.com/gh_mirrors/storm26/storm
点击查看免费下载
上一篇:3步上手pysnowball:用Python快速搭建A股行情监控小工具
下一篇:想存下抖音高清无水印视频却无从下手?douyin-downloader 这个开源工具帮你一次搞定

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询