Apache Storm Stream API 实战指南:用类型化 DSL 构建流式计算拓扑
后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载本指南围绕 Apache Storm 的 Stream APIorg.apache.storm.streams包展开介绍如何以类型化、函数式的方式表达流式计算——包括过滤、映射、窗口、聚合、连接、状态等常见操作并最终构建出可提交的 Storm 拓扑。读完本文你将掌握StreamBuilder、Stream、PairStream的核心用法能够用 Stream API 写出带窗口词频统计、按键聚合、状态查询等能力的实战拓扑。Stream API 是构建在 Storm 原生 Spout 与 Bolt 之上的一层类型化 DSL。历史上 Storm 用 Spout 和 Bolt 表达流式计算虽然简单但缺乏可复用的构造来表达过滤、变换、窗口、连接、聚合等通用流式操作Stream API 正是为弥补这一缺口而设计支持 map-reduce 风格的功能式操作。核心概念Stream、StreamBuilder 与 ValueMapperStream 与两类操作概念上一个Stream可以看作流经一条处理管道pipeline的消息流。它既可以由读取数据源如 spout产生也可以由其他 Stream 变换而来。例如// imports import org.apache.storm.streams.Stream; import org.apache.storm.streams.StreamBuilder; ... StreamBuilder builder new StreamBuilder(); // 从 spout 源得到句子流 StreamString sentences builder.newStream(new RandomSentenceSpout()).map(tuple - tuple.getString(0)); // 通过变换切分句子流得到单词流 StreamString words sentences.flatMap(s - Arrays.asList(s.split( ))); // 输出操作把单词打印到控制台 words.forEach(w - System.out.println(w));大多数流操作接受描述用户自定义行为的参数通常以 lambda 表达式形式给出例如上面的s - Arrays.asList(s.split( ))。Stream支持两类操作变换Transformations由当前流产生另一个流如上面的flatMap输出操作Output operations产生一个结果如上面的forEach。StreamBuilderStreamBuilder提供创建新流的 builder 接口通常由一个 spout 作为流的源头StreamBuilder builder new StreamBuilder(); StreamTuple sentences builder.newStream(new TestSentenceSpout());StreamBuilder会跟踪通过 Stream 表达的整体操作管道。之后可以通过build()创建 Storm 拓扑并像普通拓扑一样用StormSubmitter提交StormSubmitter.submitTopologyWithProgressBar(test, new Config(), streamBuilder.build());从源码结构看StreamBuilder.java 内部维护了一张有向图JGraphT 的DefaultDirectedGraphNode, EdgenewStream(IRichSpout)/newStream(IRichSpout, int parallelism)会创建SpoutNode并为其分配唯一的 spout id 与并行度build()则通过TopologicalOrderIterator遍历图中的节点把连续的ProcessorNode分组并映射为实际的 Bolt普通 Bolt、窗口 Bolt 或状态化 Bolt最终交给TopologyBuilder生成StormTopology。因此Stream API 的管道在构建期会被翻译成标准 Storm 拓扑组件。ValueMapper把 Tuple 提取为类型化值Value mapper 用于从 spout 发出的 tuple 中提取指定字段从而得到类型化的值流作为参数传给StreamBuilder.newStreamStreamBuilder builder new StreamBuilder(); // 提取 tuple 的第一个字段得到 StreamString 的句子流 StreamString sentences builder.newStream(new TestWordSpout(), new ValueMapperString(0));Storm 通过Pair和Tuple类Tuple3到Tuple10提供强类型 tuple可以用TupleValueMapper得到类型化 tuple 流// 提取 spout 发出 tuple 的前三个字段产生类型化 tuple 流 StreamTuple3String, Integer, Long stream builder.newStream(new TestSpout(), TupleValueMappers.of(0, 1, 2));从实现上看TupleValueMapper.java 是FunctionTuple, T的子接口而StreamBuilder中newStream(spout, valueMapper)的实现本质上是newStream(spout).map(valueMapper)——即先拿到 Tuple 流再用 mapper 做一次映射变换这是理解值流如何从 tuple 流派生出来的关键实现细节。基础变换filter、map 与 flatMapStream API 的全部接口定义在 Stream.java 和 PairStream.java 中目前支持变换、过滤、窗口、聚合、分支、连接、状态、输出与调试等大量操作。filterfilter返回由满足给定Predicate谓词返回 true的元素组成的流StreamString logs ... StreamString errors logs.filter(line - line.contains(ERROR));上例把包含 ERROR 的日志行过滤到错误流中便于后续进一步处理。mapmap返回将给定映射函数应用到流中每个值后得到的结果流。注意结果流的类型可以与原流不同StreamString words ... StreamInteger wordLengths words.map(String::length);上例通过对每个单词应用String.length从单词流得到单词长度流。flatMapflatMap返回把每个值替换为映射函数产生的内容0 个或多个值后展平得到的流。它类似map但每个值可以映射为多个值StreamString sentences ... StreamString words sentences.flatMap(s - Arrays.asList(s.split( )));上例中lambda 把流中每个值切分为单词列表flatMap再将其展平为单词流。窗口操作Windowingwindow操作产生一个窗口流包含落在窗口参数指定范围内的元素。底层窗口 Bolt 支持的所有窗口选项都能通过 Stream API 使用StreamT windowedStream stream.window(Window?, ? windowConfig);windowConfig参数指定窗口配置如基于时间或事件数量的滑动/滚动窗口// 基于时间的滑动窗口 stream.window(SlidingWindows.of(Duration.minutes(10), Duration.minutes(1))); // 基于数量的滑动窗口 stream.window(SlidingWindows.of(Count.of(10), Count.of(2))); // 滚动窗口 stream.window(TumblingWindows.of(Duration.seconds(10))); // 指定事件时间戳字段并声明迟到 tuple 流 stream.window(TumblingWindows.of(Duration.seconds(10) .withTimestampField(ts) .withLateTupleStream(late_events)));窗口操作把连续的值流切分为子集是执行连接Joins与聚合Aggregations的前提。补充一点源码事实SlidingWindows.java 提供了四种of重载分别支持Count/Duration在窗口长度与滑动间隔上的任意组合此外还支持withTimestampField指定 tuple 中携带事件时间戳的字段名、withLateTupleStream指定迟到 tuple 的输出流以及withLag(Duration)限制 tuple 时间戳的最大乱序程度。这些配置最终会映射到底层BaseWindowedBolt对应的窗口语义。变换为键值对mapToPair 与 flatMapToPair这两个操作把值流变换为键值对流StreamInteger integers … // 1, 2, 3, 4, ... PairStreamInteger, Integer squares integers.mapToPair(x - Pair.of(x, x*x)); // (1, 1), (2, 4), (3, 9), (4, 16), ...键值对流是执行groupByKey、aggregateByKey、连接等操作的前提。flatMapToPair与之类似但每个值可以产生 0 个或多个键值对。在实现层面mapToPair/flatMapToPair返回的PairStreamK, V会把输出字段声明为key、value两个字段为后续按 key 分区与聚合做好准备。聚合操作Aggregations聚合操作对流中的值或键值进行汇总。通常聚合操作作用于窗口流在每次窗口激活时发出聚合结果。aggregate 与 reduceaggregate和reduce计算全局聚合即所有分区的值被转发到单一 task 进行计算StreamLong numbers … // 聚合数字产生最近 10 秒的和 StreamLong sums numbers.window(TumblingWindows.of(Duration.seconds(10))).aggregate(new Sum()); // 用 reduce 计算最近 10 秒的和 StreamLong sums numbers.window(...).reduce((x, y) - x y);aggregate与reduce的计算方式不同reduce反复应用给定的 reducer把两个值归约为一个值直到只剩一个值。这种模式并非对所有聚合都适用或容易实现例如求平均值。aggregate做可变归约mutable reduction在处理值的过程中把结果累积到 accumulator 中。聚合操作aggregate和reduce会尽可能先在本地做分区内聚合再进行网络 shuffle以最小化网络上传输的消息量。例如计算 sum 时先计算每个分区的部分和只把部分和通过网络传到目标 bolt在目标 bolt 上合并部分和得到最终结果。为此aggregate使用CombinerAggregator接口作为参数。例如上例中的Sum可以这样实现为CombinerAggregatorpublic class Sum implements CombinerAggregatorLong, Long, Long { // sum 的初始值 Override public Long init() { return 0L; } // 把值可能是部分和累加到聚合结果上 Override public Long apply(Long aggregate, Long value) { return aggregate value; } // 合并部分和 Override public Long merge(Long accum1, Long accum2) { return accum1 accum2; } // 从 accumulator 提取结果这里 accumulator 与结果相同 Override public Long result(Long accum) { return accum; } }从 Stream.java 的实现看先本地聚合、再全局合并的优化路径由combine方法体现当流当前并行度大于 1shouldPartition()时会先插入AggregateProcessor做分区内聚合再通过全局分区global()把部分结果汇聚到一个 task 上最后用MergeAggregateProcessor合并出最终结果当并行度为 1 时则直接做分区聚合避免不必要的网络传输。aggregateByKey 与 reduceByKey这两个操作与aggregate/reduce类似但按 key 分别聚合。aggregateByKey使用给定 Aggregator 对流的每个 key 聚合值StreamString words ... // 窗口化的单词流 PairStreamString, Long wordCounts words.mapToPair(w - Pair.of(w, 1)) // 转为 (word, 1) 键值对流 .aggregateByKey(new Count()); // 计算每个单词的计数reduceByKey对流的每个 key 反复应用 reducer对值做归约StreamString words ... // 窗口化的单词流 PairStreamString, Long wordCounts words.mapToPair(w - Pair.of(w, 1)) // 转为 (word, 1) 键值对流 .reduceByKey((x, y) - x y); // 计算每个单词的计数与全局聚合一样按 key 的分区内局部聚合也会先计算部分结果发送到目标 bolt 后再合并出最终聚合结果。groupByKeygroupByKey对键值对流返回一个新流其中的值按 key 分组// 一个 (user, score) 键值对流如 (alice, 10), (bob, 15), (bob, 20), (alice, 11), (alice, 13) PairStreamString, Double scores ... // 最近窗口内每个用户的分数列表如 (alice, [10, 11, 13]), (bob, [15, 20]) PairStreamString, IterableInteger userScores scores.window(...).groupByKey();countByKeycountByKey对流中每个 key 的值计数StreamString words ... // 窗口化的单词流 PairStreamString, Long wordCounts words.mapToPair(w - Pair.of(w, 1)) // 转为 (word, 1) 键值对流 .countByKey(); // 计算每个单词的计数从文档说明看countByKey内部使用aggregateByKey来完成计数仓库内置的Count聚合器即实现于此。重新分区Repartitionrepartition操作对当前流重新分区返回具有指定分区数的新流后续操作将在该并行度上执行可用于提升或降低流中操作的并行度。初始分区数也可以在创建流时通过StreamBuilder.newStream指定// 流 s1 有 2 个分区s1 上的操作在该并行度执行 StreamString s1 builder.newStream(new TestWordSpout(), new ValueMapperString(0), 2); // 流 s2 及后续操作将有 3 个分区 PairStreamString, Integer s2 s1.map(function1).repartition(3); // 在 s2 上执行 map 并打印结果 s2.map(function2).print();注意repartition意味着网络传输。上例中第一个 map 操作function1会在 2 个分区s1 的 2 个分区上执行第二个 map 操作function2会在 3 个分区s2 的 3 个分区上执行这同时意味着两个 map 操作必须分别在两个 bolt 上执行并产生网络传输。从源码看Stream.java 的repartition(int parallelism)会校验并行度必须 1若新并行度与当前节点一致则直接返回自身不产生多余的分区节点否则插入一个PartitionNode作为新的流节点。输出操作Output operations输出操作把流中变换后的值推送到控制台、外部存储数据库、文件甚至 Storm bolt。printprint把流中的值打印到控制台// 把单词转大写并打印 words.map(String::toUpperCase).print();peekpeek返回由流中元素组成的流同时在元素从结果流中消费时额外执行给定动作可用于观察流中任何阶段的流动值builder.newStream(...).flatMap(s - Arrays.asList(s.split( ))) // 值流过时打印 flatMap 的结果 .peek(s - System.out.println(s)) .mapToPair(w - new Pair(w, 1))forEachforEach是最通用的输出操作可为流中每个值执行任意代码例如把结果写入外部数据库、文件等stream.forEach(value - { // 记录日志 LOG.debug(value); // 把值存入数据库等 statement.executeUpdate(..); } );toto允许把已有 bolt 作为 sink 接入流// redisBolt 是一个标准的 storm bolt IRichBolt redisBolt new RedisStoreBolt(poolConfig, storeMapper); ... // 生成词频统计并把结果通过 redis bolt 存入 redis builder.newStream(new TestWordSpout(), new ValueMapperString(0)) .mapToPair(w - Pair.of(w, 1)) .countByKey() // (word, count) 键值对被转发给 redisBolt 存储 .to(redisBolt);注意to提供的可靠性保证只取决于该 bolt 自身提供的保证。从源码看Stream.to(IRichBolt bolt, int parallelism)会把 bolt 包装成SinkNode挂到流图上默认并行度为 1构建期由StreamBuilder.build()中的addSink完成接线与 grouping 声明。分支操作Branchbranch操作用于在流上表达 if-then-else 逻辑StreamT[] streams stream.branch(PredicateT... predicates)谓词按给定顺序依次应用到流中的值上值被转发到第一个匹配谓词对应的按下标索引的结果流若所有谓词都不匹配该值被丢弃。例如StreamInteger[] streams builder.newStream(new RandomIntegerSpout(), new ValueMapperInteger(0)) .branch(x - (x % 2) 0, x - (x % 2) 1); StreamInteger evenNumbers streams[0]; StreamInteger oddNumbers streams[1];从实现上看branch在流上插入一个BranchProcessor为每个谓词创建独立的子节点与分支流并把谓词 - 分支流的映射注册到该 processor 中从而把不同分支的值路由到不同输出流。连接Joinsjoin操作把一个流的值与另一个流中具有相同 key 的值连接起来PairStreamLong, Long squares … // (1, 1), (2, 4), (3, 9) ... PairStreamLong, Long cubes … // (1, 1), (2, 8), (3, 27) ... // 连接 squares 与 cubes 流产生 (1, [1, 1]), (2, [4, 8]), (3, [9, 27]) ... PairStreamLong, PairLong, Long joined squares.window(TumblingWindows.of(Duration.seconds(5))).join(cubes);连接通常作用于窗口流把当前窗口中每条流到达的键值按 key 连接。调用连接的那条流的并行度会延续到连接后的流上。可选参数ValueJoiner用于指定如何连接每个匹配 key 的两个值默认行为是返回来自两条流的值的Pair。左连接、右连接和全外连接均受支持。coGroupByKeycoGroupByKey把此流的值与另一条流中相同 key 的值分组到一起// 一个 (key, value) 流如 (k1, v1), (k2, v2), (k2, v3) PairStreamString, String stream1 ... // 另一个 (key, value) 流如 (k1, x1), (k1, x2), (k3, x3) PairStreamString, String stream2 ... // 最近窗口内按 key 分组后的值如 (k1, ([v1], [x1, x2]), (k2, ([v2, v3], [])), (k3, ([], [x3])) PairStreamString, IterableString coGroupedStream stream1.window(...).coGroupByKey(stream2);状态操作StateStorm 提供了保存、更新计算状态以及查询状态的 API。updateStateByKeyupdateStateByKey通过把给定的状态更新函数应用到该 key 的旧状态与新值上更新状态。它可以传入状态的初始值与状态更新函数来调用也可以直接提供一个StateUpdater实现PairStreamString, Long wordCounts ... // 更新状态中的单词计数第一个参数 0L 是状态的初始值 // 第二个参数是一个把计数加到当前状态值上的函数。 StreamStateString, Long streamState wordCounts.updateStateByKey(0L, (state, count) - state count); streamState.toPairStream().print();状态值可以是任意类型上例中是Long类型保存单词计数。内部实现上Storm 使用状态化 Bolt 存储状态可用 Storm 配置topology.state.provider选择状态提供者实现。例如设置为org.apache.storm.redis.state.RedisKeyValueStateProvider时使用基于 Redis 的状态存储对应仓库中的 external/storm-redis 模块。stateQuerystateQuery用于查询由updateStateByKey更新的状态。updateStateByKey返回的StreamState必须用于查询流状态流中的值被用作查询状态的 key// QuerySpout 发出的单词流用作查询状态的 key builder.newStream(new QuerySpout(), new ValueMapperString(0)) // 查询状态并发出匹配的 (key, value) 作为结果。 // updateStateByKey 返回的流状态作为 stateQuery 的参数传入。 .stateQuery(streamState).print();从源码看Stream.java 中stateQuery的实现会先按 key 字段做 field grouping 分区need field grouping for state query so that the query is routed to the correct task再把StateQueryProcessor挂到流上而StreamBuilder在构建时会尽量让StateQueryProcessor与对应的UpdateStateByKeyProcessor映射到同一个StatefulProcessorBolt上从而保证状态查询被路由到正确的 task。可靠性保证Guarantees目前使用 Stream API 构建的拓扑提供**至少一次at-least once**保证。需要注意当前只有updateStateByKey操作执行在底层的StatefulBolt上其他有状态操作join、windowing、aggregation 等执行在IRichBolt上并把状态保存在内存中依赖 Storm 的 ack 与重放机制来重建状态。按照该文档的说明未来 Stream API 的底层框架会增强以提供**恰好一次exactly once**保证。完整示例窗口词频统计拓扑下面是一个用 Stream API 表达的词频统计拓扑StreamBuilder builder new StreamBuilder(); builder // 两个分区的随机句子流 .newStream(new RandomSentenceSpout(), new ValueMapperString(0), 2) // 两秒滚动窗口 .window(TumblingWindows.of(Duration.seconds(2))) // 把句子切分为单词 .flatMap(s - Arrays.asList(s.split( ))) // 创建 (word, 1) 键值对流 .mapToPair(w - Pair.of(w, 1)) // 计算最近两秒窗口内的单词计数 .countByKey() // 打印结果到 stdout .print();RandomSentenceSpout是普通的 Storm spout持续发出随机句子。句子流被切分为两秒的窗口计算每个窗口内的单词计数并打印。该流可以像普通拓扑一样提交Config config new Config(); config.setNumWorkers(1); StormSubmitter.submitTopologyWithProgressBar(topology-name, config, builder.build());仓库的 storm-starter 模块提供了完整的可运行示例位于 examples/storm-starter/src/jvm/org/apache/storm/starter/streams包括WindowedWordCount.java与上文示例对应的窗口词频统计并额外加入了filter(x - x.getSecond() 5)过滤出现次数较少的单词AggregateExample.java、BranchExample.java、JoinExample.java、GroupByKeyAndWindowExample.java、StatefulWordCount.java、StateQueryExample.java、TypedTupleExample.java、WordCountToBolt.java分别演示聚合、分支、连接、按 key 分组、状态化词频、状态查询、类型化 tuple 以及把结果输出到既有 bolt 等能力。实践要点小结从 spout 起步builder.newStream(spout, valueMapper, parallelism)是创建类型化流的统一入口valueMapper 决定从 tuple 中提取哪些字段并行度决定后续操作的执行级别。聚合自动本地化aggregate/reduce以及按键版本会自动做分区内局部聚合再网络合并因此优先使用CombinerAggregator描述聚合可显著减少跨网络传输的数据量。窗口是连接与聚合的前提SlidingWindows/TumblingWindows支持时间、数量及其混合组合事件时间处理需要配合withTimestampField与withLateTupleStream。状态能力有限定updateStateByKey运行在状态化 Bolt 上并可用topology.state.provider切换存储后端其他有状态操作当前为内存态、依赖 ack/重放整体拓扑为至少一次语义。repartition 意味着网络传输调整并行度会引入新的分区边界连续的高并行度操作会被拆分到不同 bolt实际编码时需权衡并行度收益与网络开销。赞分享后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载相关推荐CANN pypto 分阶段写回枚举 STPhase 详解与 AccPhase 的 unit_flag 握手机制与实战用法CANN pypto 分阶段写回枚举 STPhase 详解与 AccPhase 的 unit_flag 握手机制与实战用法 分阶段写回store phase大数据流处理后端Apache Pulsar 与 Apache Storm 集成实战基于 Pulsar Spout / Bolt 适配器打通流式计算拓扑Apache Pulsar 与 Apache Storm 集成实战基于 Pulsar Spout / Bolt 适配器打通流式计算拓扑 Apache Puls消息队列后端流处理为 Apache Storm 定义非 JVM 语言 DSL基于 Thrift 结构与 storm shell 的多语言拓扑构建指南为 Apache Storm 定义非 JVM 语言 DSL基于 Thrift 结构与 storm shell 的多语言拓扑构建指南 导读 Apache Sto流处理后端大数据上一篇如何高效处理go-swagger枚举类型从Swagger规范到Go常量的完整指南下一篇OpenRig 学徒交接中的 Orchestrator 角色判断边界、权威门控与交接前 READY 验证创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考