资讯动态

Apache Storm Clojure DSL 完全指南:用 Clojure 定义 Spout、Bolt 与 Topology

发布时间:2026/10/9 2:08:43 来源:尧图企业网站定制
后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载本文是 Apache Storm 官方文档《Clojure DSL》docs/Clojure-DSL.md的深度扩展版本面向希望用纯 Clojure 编写 Storm 流处理作业的开发者。读完本文你将掌握topology/spout-spec/bolt-spec的拓扑组装方式defbolt/defspout的三种形态简单、参数化、预处理多语言 shell bolt 的接入方法以及如何在本地集群与真实集群上提交、测试拓扑——全部只写 Clojure不触碰 Java。1. 认识 storm-clojureClojure DSL 的定位与入口Storm 的核心 API 是 Java 的IRichSpout、IRichBolt、TopologyBuilder等但通过独立的storm-clojure模块官方为 Clojure 用户提供了完整的 DSLspout、bolt、拓扑的组装与提交都可以用 Clojure 完成。正如文档开头所述Clojure DSL 能访问 Java API 暴露的一切能力因此 Clojure 开发者可以不写一行 Java 就完成 Storm 拓扑的开发。DSL 的实现位于 storm-clojure/src/clj/org/apache/storm/clojure.clj 命名空间org.apache.storm.clojure其中定义了defbolt、defspout、bolt、spout、emit-bolt!、emit-spout!等核心宏与函数。而拓扑组装层的topology、bolt-spec、spout-spec、shell-bolt-spec则是对 storm-clojure/src/clj/org/apache/storm/thrift.clj 中 Thrift 构造函数的别名defalias见 clojure.clj 第 216-220 行负责把 DSL 声明翻译成 Storm 内部的 Thrift 拓扑结构。本文将从定义拓扑到测试拓扑依次展开覆盖 DSL 的全部五个组成部分。2. 定义拓扑topology、spout-spec与bolt-spec2.1topology函数与组件 ID使用topology函数定义拓扑它接收两个参数spout spec 的 map和bolt spec 的 map。每个 spec 通过声明输入、并行度等信息把组件代码接线进拓扑。官方文档给出的经典示例来自 storm-starter 示例工程 examples/storm-starter/src/clj/org/apache/storm/starter/clj/word_count.clj(topology {1 (spout-spec sentence-spout) 2 (spout-spec (sentence-spout-parameterized [the cat jumped over the door greetings from a faraway land]) :p 2)} {3 (bolt-spec {1 :shuffle 2 :shuffle} split-sentence :p 5) 4 (bolt-spec {3 [word]} word-count :p 6)})两个 map 均以**组件 IDcomponent id**为键、以 spec 为值。与 Java 方式一致组件 ID 必须全局唯一跨 spout 与 bolt 的 map 也不能重复并且在声明 bolt 输入时通过组件 ID 来引用上游组件。从源码看topology底层就是 thrift.clj 的 mk-topology它遍历 spout/bolt map依次调用TopologyBuilder.setSpout/setBolt把:p作为并行度 hint、:conf作为组件配置写入最后createTopology生成StormTopology。2.2spout-spec指定 Spout 与并行度spout-spec接收两个参数spout 实现对象实现IRichSpout接口的对象例如用defspout定义的 spout或TestWordSpout这类 Java 类实例以及可选的关键字参数。目前唯一的选项是:p用于指定该 spout 的并行度task 数省略:p时 spout 以单 task 运行(spout-spec (sentence-spout-parameterized [the cat jumped over the door greetings from a faraway land]) :p 2)在 thrift.clj 的 mk-spout-spec 中可以看到mk-spout-spec会返回{:obj spout :p parallelism-hint :conf conf}这样一个 maptopology再从:p中取出并行度传给TopologyBuilder。值得注意:p与:parallelism-hint是等价的别名二者同时存在时:p优先。2.3bolt-spec输入声明与流分组bolt-spec接收三个参数输入声明input declaration、bolt 实现对象实现IRichBolt以及可选关键字参数。输入声明是一个从流 ID 到流分组stream grouping的 map。流 ID 有两种形式[组件ID 流ID]订阅该组件上的指定命名流组件ID裸字符串订阅该组件的默认流default stream。流分组grouping可取以下关键字分组关键字含义:shuffleshuffle 分组随机分发字段名向量如[id name]fields 分组按指定字段哈希分发:globalglobal 分组全部送往 task 0:allall 分组广播给所有 task:directdirect 分组由上游显式指定目标 task官方文档给出的多流输入声明示例{[2 1] :shuffle 3 [field1 field2] [4 2] :global}它订阅了三条流组件 2 的流 1shuffle 分组、组件 3 的默认流按 field1、field2 字段分组、组件 4 的流 2global 分组。在源码层thrift.clj 的 mk-inputs 与 mk-grouping 负责把这套 Clojure 声明翻译为 ThriftGlobalStreamIdGrouping流 ID 是向量时用[组件ID 流ID]构造GlobalStreamId是裸字符串时补上默认流 ID分组关键字依次映射到Grouping/shuffle、Grouping/fields、Grouping/global、Grouping/all、Grouping/direct。从代码可以推断mk-grouping还额外支持:local-or-shuffle、:none、自定义CustomStreamGrouping对象甚至JavaObject这比文档列出的五种关键字更丰富。与spout-spec相同bolt-spec目前唯一支持的选项也是:p并行度。流分组的更多概念说明可参考 docs/Concepts.md。2.4shell-bolt-spec接入非 JVM 语言的 Boltshell-bolt-spec用于定义由非 JVM 语言实现的 bolt。它接收的参数依次为输入声明、要运行的命令行程序、实现 bolt 的脚本文件名、输出声明以及bolt-spec所支持的关键字参数。官方文档给出的示例(shell-bolt-spec {1 :shuffle 2 [id]} python3 mybolt.py [outfield1 outfield2] :p 25)该声明创建了一个并行度为 25 的 Python bolt订阅组件 1 的默认流shuffle 分组与组件 2 的默认流按 id 字段分组用python3 mybolt.py进程实现输出字段为 outfield1、outfield2。从 thrift.clj 的 mk-shell-bolt-spec 与 shell-component-params 可以看出两种调用约定若第二参数是字符串脚本路径则被解释为单条命令[command script]若第二参数本身是序列则整个被当作命令数组可携带参数并自动把RichShellBolt包装进mk-bolt-spec。多语言multilang子进程协议、stdin/stdout 通信的完整细节参见 docs/Using-non-JVM-languages-with-Storm.md。输出声明的语法与defbolt一致见 3.5 节。3.defbolt在 Clojure 中定义 Bolt3.1 为什么不能直接 reifyIRichBoltBolt 必须可序列化——因为 Storm 会把组件代码随拓扑提交到集群的各个 worker。Clojure 的闭包closure不可序列化所以直接reifyIRichBolt并不可行。defbolt正是为了解决这一约束而存在它把 bolt 的实现函数与参数分离在宏展开时生成可序列化的ClojureBolt包装对象详见 3.2 的源码说明同时提供了比手写 Java 接口更简洁的语法。defbolt的完整签名(defbolt _name_ _output-declaration_ *_option-map_ _impl_)省略选项 map 等价于{:prepare false}。三种形态简单、参数化、预处理分别对应不同的使用场景。3.2 简单 Bolt直接实现execute省略选项 map 时定义的是非预处理non-preparedboltDSL 只要求提供IRichBolt.execute方法的实现。实现接收两个参数——tuple 与OutputCollector后跟execute的函数体。DSL 会自动为参数添加类型提示type-hint因此使用 Java interop如(.getString tuple 0)时无需担心反射开销。官方文档的句子切分示例(defbolt split-sentence [word] [tuple collector] (let [words (.split (.getString tuple 0) )] (doseq [w words] (emit-bolt! collector [w] :anchor tuple)) (ack! collector tuple) ))定义完成后split-sentence被绑定到一个真正的IRichBolt对象可以直接用于拓扑(bolt-spec {1 :shuffle} split-sentence :p 5)从 clojure.clj 中 defbolt 宏的展开逻辑可以看到非预处理形态下宏会把[tuple collector]拆解——collector 被提升到闭包中execute只接收 tuple 一个参数并用bolt宏reifyIBolt包装。而ClojureBoltClojureBolt.java在prepare阶段通过ClojureUtil.loadClojureFn加载 var、应用参数、调用实现函数得到真正的IBoltexecute阶段再把输入包成ClojureTuple转发——这就是Clojure 函数可被序列化地塞进拓扑的底层机制。3.3 参数化 Bolt:params选项很多场景需要给 bolt 传运行时参数。在选项 map 中放入:params即可例如给每个输入字符串追加后缀的 bolt(defbolt suffix-appender [word] {:params [suffix]} [tuple collector] (emit-bolt! collector [(str (.getString tuple 0) suffix)] :anchor tuple) )与简单 bolt 不同**指定:params后suffix-appender被绑定为一个返回IRichBolt的函数**而非对象本身。使用时需要先调用该函数(bolt-spec {1 :shuffle} (suffix-appender -suffix) :p 10)从宏实现看:params会让defbolt生成defn接收可变参数并通过clojure-bolt把参数原样传给ClojureBolt的构造器最终在prepare时applyTo到实现函数上——参数也因此被序列化进拓扑。3.4 预处理 Boltprepared bolt:prepare true与bolt宏需要跨 tuple 维护状态如 join、流式聚合时用:prepare true定义预处理 bolt。其实现是一个接收拓扑配置conf、TopologyContext、OutputCollector并返回IBolt实现的函数。这种设计允许在execute、cleanup周围形成闭包。官方文档的词频统计示例与 word_count.clj 完全一致(defbolt word-count [word count] {:prepare true} [conf context collector] (let [counts (atom {})] (bolt (execute [tuple] (let [word (.getString tuple 0)] (swap! counts (partial merge-with ) {word 1}) (emit-bolt! collector [word (counts word)] :anchor tuple) (ack! collector tuple) )))))这里词频被保存在闭包中的counts一个 atom 包着的 map。bolt宏是比 reify 更简洁、且自动加类型提示的IBolt实现方式clojure.clj 的 bolt 宏内部实际就是 reifyIBolt。注意预处理 bolt 的execute只接收 tuple——因为OutputCollector已在闭包中简单 bolt 的 execute 才有 collector 第二参数。预处理 bolt 同样可以参数化:params与:prepare true可同时使用。storm-starter 的滚动 Top-N 示例 examples/storm-starter/src/clj/org/apache/storm/starter/clj/bolts.clj 是这两种能力组合的典型rolling-count-bolt同时使用:prepare true、:params [window-length emit-frequency]与:conf {TOPOLOGY-TICK-TUPLE-FREQ-SECS emit-frequency}把窗口长度、发射频率作为参数注入并利用 tick tuple 定时驱动窗口滑动——这展示了defbolt在真实生产级组件中的用法。3.5 输出声明output declaration语法Clojure DSL 用简洁语法声明 bolt 输出。最一般的形式是从流 ID 到流 spec 的 map{1 [field1 field2] 2 (direct-stream [f1 f2 f3]) 3 [f1]}流 ID 是字符串流 spec 是字段向量或用direct-stream包裹的字段向量——后者将该流标记为 direct streamdirect 流只能通过emit-direct-bolt!/emit-direct-spout!定向发射概念详见 docs/Concepts.md。若 bolt 只有一条输出流可省略 map直接用向量声明默认流的字段[word count]这等价于声明默认流上字段为[word count]的输出。底层处理见 thrift.clj 的 mk-output-spec向量会被自动包成{Utils/DEFAULT_STREAM_ID 向量}并转成 ThriftStreamInfodirect-stream即(StreamInfo. fields true)置 direct 标志。3.6 发射、确认与失败emit-bolt!/emit-direct-bolt!/ack!/fail!与其直接调用OutputCollector的 Java 方法DSL 提供了更顺手的封装函数定义见 clojure.cljemit-bolt!参数为OutputCollector、要发射的值序列Clojure 序列关键字参数:anchor与:stream。:anchor可以是单个 tuple 或 tuple 列表:stream是目标流 ID。省略关键字参数时向默认流发射未锚定unanchoredtuple。emit-direct-bolt!参数为OutputCollector、目标 task ID、值序列关键字参数:anchor与:stream。只能向声明为 direct stream 的流发射。ack!参数为OutputCollector与要确认的 tuple。fail!参数为OutputCollector与要失败的 tuple。从实现看emit-bolt!最终调用OutputCollector.emit(stream, anchor, values)且:anchor会经过collectify统一成列表ack!/fail!直接委托.ack/.fail。此外源码中还提供了reset-timeout!重置 tuple 超时与report-error!上报错误两个未在文档中列出的辅助函数需要时同样可用。锚定与确认机制的完整原理参见 docs/Guaranteeing-message-processing.md。4.defspout在 Clojure 中定义 Spout4.1 基本 spout 与可靠性语义与 bolt 同理spout 也必须可序列化不能直接 reifyIRichSpout。defspout提供同样的规避机制与简洁语法签名与defbolt相同(defspout _name_ _output-declaration_ *_option-map_ _impl_)与defbolt相反defspout省略选项 map 时默认{:prepare true}。输出声明语法与defbolt完全一致。官方文档的sentence-spout示例同样来自 word_count.clj(defspout sentence-spout [sentence] [conf context collector] (let [sentences [a little brown dog the man petted the dog four score and seven years ago an apple a day keeps the doctor away]] (spout (nextTuple [] (Thread/sleep 100) (emit-spout! collector [(rand-nth sentences)]) ) (ack [id] ;; You only need to define this method for reliable spouts ;; (such as one that reads off of a queue like Kestrel) ;; This is an unreliable spout, so it does nothing here ))))实现函数接收拓扑配置、TopologyContext与SpoutOutputCollector返回一个ISpout对象。spout宏与bolt宏对应是 reifyISpout的简洁写法clojure.clj 的 spout 宏。可靠性语义这是一个不可靠unreliablespout——发射 tuple 时不带消息 ID因此ack/fail永远不会被调用。可靠 spout需要在发射时提供消息 IDtuple 处理完成或失败后ack/fail回调才会被触发。可靠性的完整机制参见 docs/Guaranteeing-message-processing.md。4.2emit-spout!与emit-direct-spout!emit-spout!的参数为SpoutOutputCollector与新 tuple 的值序列关键字参数:stream目标流与:id消息 ID供ack/fail回调使用。省略时向默认流发射无消息 ID 的 tuple。emit-direct-spout!则在第二参数位置多一个目标 task ID用于向 direct stream 定向发射实现见 clojure.clj。4.3 参数化与 unprepared spoutspout 也可以参数化——此时符号绑定为返回IRichSpout的函数。还可以声明unprepared spout:prepare false只需定义nextTuple方法。官方文档的运行时参数化 spout 示例(defspout sentence-spout-parameterized [word] {:params [sentences] :prepare false} [collector] (Thread/sleep 500) (emit-spout! collector [(rand-nth sentences)]))在拓扑中使用(spout-spec (sentence-spout-parameterized [the cat jumped over the door greetings from a faraway land]) :p 2)从 defspout 宏的展开逻辑看:prepare false时collector 被提升进闭包、只生成nextTuple的实现并通过spout宏包成ISpout:params存在时生成defn返回由ClojureSpout包装的对象。5. 提交拓扑本地模式与集群模式5.1 使用StormSubmitter拓扑定义完成后无论本地模式还是远程集群模式提交方式都与 Java 一致——直接使用StormSubmitter类。官方推荐的提交入口见 word_count.clj 与 exclamation.clj 的submit-topology!(StormSubmitter/submitTopology name {TOPOLOGY-DEBUG true TOPOLOGY-WORKERS 3} (mk-topology))DSL 还封装了便捷函数submit-remote-topologyclojure.clj等价于调用StormSubmitter/submitTopology。5.2org.apache.storm.config命名空间所有配置常量创建拓扑配置最省事的方式是使用org.apache.storm.config命名空间storm-clojure/src/clj/org/apache/storm/config.clj。它遍历 JavaConfig类的全部声明字段为每个配置项自动生成一个 Clojure 常量——常量名与Config类的静态常量一一对应只是把下划线换成连字符。例如Config.TOPOLOGY_WORKERS对应TOPOLOGY-WORKERSConfig.TOPOLOGY_DEBUG对应TOPOLOGY-DEBUG。代码中可见每个常量都直接.引用 JavaConfig字段的值。官方文档给出的配置示例——15 个 worker 且开启 debug 模式{TOPOLOGY-DEBUG true TOPOLOGY-WORKERS 15}由于是自动生成的TOPOLOGY-*、NIMBUS-*、SUPERVISOR-*等数百个常量都可以直接通过连字符命名访问。完整配置项说明可对照 docs/Configuration.md 与 conf/storm.yaml.example。5.3 本地集群环境DSL 还提供了local-cluster函数clojure.clj通过延迟求值实例化org.apache.storm.LocalCluster便于在开发环境快速跑通拓扑源码注释特别指出该延迟求值是为了避免LocalCluster - testing - nimbus - bootstrap - clojure - LocalCluster的循环依赖。6. 测试拓扑Storm 为 Clojure 提供了强大的内置测试设施集中在独立的storm-clojure-test模块核心工具位于 storm-clojure-test/src/clj/org/apache/storm/testing.cljwith-local-clustertesting.clj 第 126 行在内存中拉起一整套模拟集群含 Nimbus、Supervisor、ZooKeeper 均可配置在宏体内提交并验证拓扑结束自动清理。with-simulated-time-local-cluster配合Time$SimulatedTime使用模拟时钟配合advance-cluster-time可以快进模拟时间是测试窗口聚合、tick tuple 类拓扑的关键手段。with-tracked-clustertracked-wait第 235、244 行Storm 的核心测试利器。with-tracked-cluster会为所有 tuple 打上跟踪标记tracked-wait阻塞直到拓扑空闲且 spout 累计发射了指定数量的 tuple——这让发射 N 条、断言收到 N 条的确定性测试成为可能。submit-local-topology/submit-local-topology-with-opts向测试 Nimbus 提交拓扑底层校验配置可 JSON 序列化后调用Nimbus.submitTopology。test-tuple第 253 行直接构造测试用Tuple可指定 stream、component、fields用于对单个 bolt 做单元级验证。官方文档在Testing topologies一节引用过外部博客文章来介绍这些设施就当前仓库而言这些能力的权威来源就是上述 testing.clj 文件本身以及 storm-clojure 与 storm-server 中大量以此为基座的测试用例。实际测试拓扑时推荐组合使用with-local-cluster或with-tracked-cluster提交拓扑、用tracked-wait等待确定性进度再对输出做断言。7. 完整实战WordCount 拓扑全解析最后把本文所有知识点串成一条完整的生产线。以下就是 examples/storm-starter/src/clj/org/apache/storm/starter/clj/word_count.clj 的完整结构该文件同时也是官方 Clojure DSL 文档的主示例① 定义两个 spout不可靠的sentence-spout随机发射句子与参数化、unprepared 的sentence-spout-parameterized运行时注入句子集合。② 定义两个 bolt简单 boltsplit-sentence切分句子为单词并 ack与预处理 boltword-count用 atom 闭包维护词频。③ 组装拓扑(topology {1 (spout-spec sentence-spout) 2 (spout-spec (sentence-spout-parameterized [the cat jumped over the door greetings from a faraway land]) :p 2)} {3 (bolt-spec {1 :shuffle 2 :shuffle} split-sentence :p 5) 4 (bolt-spec {3 [word]} word-count :p 6)})数据流为spout 1/2并行度 2→ bolt 3shuffle 分组、并行度 5→ bolt 4按 word 字段分组、并行度 6词频结果从 bolt 4 的默认流输出。④ 提交-main调用submit-topology!通过StormSubmitter/submitTopology提交{TOPOLOGY-DEBUG true TOPOLOGY-WORKERS 3}配置的拓扑。若在本地开发把StormSubmitter换成local-clustersubmit-local-topology即可在内存集群中验证。另可参考 exclamation.clj同一 spout 扇出到两串 bolt、字符串拼接 !!! 的极简示例与 bolts.cljrolling-count-bolt等组合:prepare true/:params/:conf的生产级组件作为进一步练习的范本。赞分享后端大数据【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm22/storm点击查看免费下载相关推荐Lance数据湖格式构建面向AI工作流的高性能存储架构终极指南Lance数据湖格式构建面向AI工作流的高性能存储架构终极指南 在当今数据驱动的AI时代传统数据湖方案正面临着前所未有的性能挑战。Lance数据湖格式作为专数据库向量数据库数据湖全文检索CANN pypto 分阶段写回枚举 STPhase 详解与 AccPhase 的 unit_flag 握手机制与实战用法CANN pypto 分阶段写回枚举 STPhase 详解与 AccPhase 的 unit_flag 握手机制与实战用法 分阶段写回store phase大数据流处理后端Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战 本篇技术指消息队列后端流处理上一篇Data-Juicer 缓存管理完全指南指纹命中、压缩存储与临时目录下一篇patent-disclosure-skill 交底书生成前摘要预览Step 6实战指南类型路由、按类型裁剪的摘要口径与确认交互创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑