
Storm这个框架这几年在实时计算圈子里热度虽然被Flink压了一头但老项目里存量巨大尤其是金融、运营商、物联网这些对实时性要求极高的场景Storm依然是不可替代的中坚力量。很多人觉得Storm难上手其实主要卡在集群搭建这一关——单机模式跑个WordCount谁都会可一旦要上生产环境涉及的角色分配、Zookeeper协调、资源隔离、高可用配置每一个环节都是坑。这篇东西我把整个搭建过程掰开揉碎讲清楚从架构原理到配置文件每一个参数的含义再到验证和常见故障排查全部覆盖。不管你是刚入门的大数据新人还是被公司派去维护老集群的运维照着一步步操作都能搭出一套稳定可用的Storm集群。1. 搭建前必须搞清楚的架构与准备工作1.1 Storm集群的三个核心角色到底在干什么Storm集群不像Hadoop那样有三五个组件十几个服务它最核心的就三个角色Nimbus、Supervisor和UI。很多人搭集群失败问题不在安装步骤而是压根没搞懂这三个角色之间的协作关系导致配置和角色分配全乱套。先看Nimbus。它是集群的“大脑”负责接收你提交的Topology计算任务然后把任务分配下去。Nimbus本身不干活它只做调度和监控——哪个Worker进程死了、哪个任务需要重新分配、哪个Supervisor节点资源不够了全靠Nimbus决策。关键点在于Nimbus的调度决策是存储在Zookeeper里的Nimbus进程本身是无状态的。这意味着Nimbus挂了已经提交的任务不会立刻中断只是集群失去调度能力。所以在高可用配置里Nimbus可以部署多个用Zookeeper选主。再看Supervisor。它是每台工作节点上的“监工”负责听从Nimbus安排在本机启动Worker进程执行具体的数据处理任务。每个Supervisor节点可以配置多个Worker槽位也就是可以同时跑多个Worker进程并行处理多个数据流。这里有个新手容易忽略的点Topology执行时Executor运行在Worker进程内部一个Worker进程可以包含多个Executor而Executor的个数直接决定任务的并行度。所以Supervisor节点的资源配置直接决定了整个集群的吞吐上限。最后是UI官方名字叫Storm UI。它是个Web服务提供集群状态的可视化界面展示Nimbus和Supervisor的存活状态、Topology运行状况、延迟指标等。虽然UI不参与数据流处理但生产环境里我建议一定部署并且配上告警——否则集群出问题你根本来不及发现。1.2 版本选型和依赖组件规划关于版本踩坑我可以直接说结论如果你的生产环境是JDK 8建议用Storm 1.2.x如果能接受JDK 11Storm 2.x也可以考虑。但网上很多教程还在写Storm 0.9.x/0.10.x那些版本太老了ZeroMQ的兼容性、Kerberos支持、UI功能都有明显缺陷新项目没必要碰。Storm的核心依赖是Zookeeper。它负责保存集群元数据、任务分配信息、Nimbus主节点选举结果。Zookeeper集群本身最好也搭成奇数节点比如3台或者5台保证选举可用。如果要做高性能的模式还需要一个消息队列配合业界最常用的就是Kafka把Kafka当作数据源把Storm当作实时处理引擎。如果你的实时链路是Kafka Storm HBase/HDFS这套组合那么Kafka版本建议与Storm的KafkaClient插件兼容1.2.x配Kafka 2.x是比较稳妥的组合。另外JDK版本一定要提前确认。Storm 1.2.x要求JDK 8不能高于1.8否则运行时会报UnsupportedClassVersionError。我在生产上就遇到过有人在JDK 11上强行跑Storm 1.1.2结果Nimbus启动后直接崩溃的情况。1.3 服务器规划与资源评估这里我给出一个典型的3节点高可用集群规划这个规划方式可以灵活扩展到几十个节点。节点角色资源配置建议核心部署组件节点A控制节点8C16GNimbus、Zookeeper、Storm UI节点B工作节点16C32GSupervisor节点C工作节点16C32GSupervisor如果你要更稳妥的高可用可以部署双Nimbus主备或者用3个节点都跑Zookeeper这样任何一个节点宕机元数据都不会丢。但要注意Nimbus和Zookeeper竞争同一台机器的资源时建议机器内存不要低于16G因为Zookeeper频繁写磁盘Nimbus需要加载JAR包和做任务调度两者都是资源消耗大户。资源评估的核心指标是每个Worker进程建议分配2-4G堆内存每台Supervisor节点建议至少分配4-8个Worker槽位。这个比例需要根据你的Topology实际并行度和线上数据量来调整后面在调优部分详细展开。2. 环境基础搭建与Zookeeper集群初始化2.1 JDK安装与环境变量配置不论你是Ubuntu还是CentOS这一步都差不多。我是CentOS 7.9的环境直接解压JDK到指定目录然后配置/etc/profile环境变量。这里有一个小细节我踩过坑不要只配置JAVA_HOMEPATH里也要带上$JAVA_HOME/bin否则后面的启动脚本会找不到Java命令。# 解压JDK tar -zxvf jdk-8u202-linux-x64.tar.gz -C /opt/ # 配置环境变量 echo export JAVA_HOME/opt/jdk1.8.0_202 export PATH$JAVA_HOME/bin:$PATH /etc/profile source /etc/profile # 验证 java -version运行后的输出应该是java version 1.8.0_202。这里强调一下JDK路径尽量固定在/opt这类标准目录下不要安装在带空格的路径里Storm脚本解析路径时会出问题。2.2 Zookeeper集群搭建与配置详解Zookeeper的配置在Storm集群中是最基础的一环。我见过不少人在Zookeeper上栽跟头配置写错一个字母Nimbus连接不上报各种session timeout异常。# 下载解压 wget https://archive.apache.org/dist/zookeeper/zookeeper-3.4.14/zookeeper-3.4.14.tar.gz tar -zxvf zookeeper-3.4.14.tar.gz -C /opt/ cp /opt/zookeeper-3.4.14/conf/zoo_sample.cfg /opt/zookeeper-3.4.14/conf/zoo.cfg然后编辑zoo.cfg核心配置如下tickTime2000 initLimit10 syncLimit5 dataDir/data/zookeeper clientPort2181 server.1node01:2888:3888 server.2node02:2888:3888 server.3node03:2888:3888这里解释一下里面的几个关键参数。tickTime是Zookeeper的基本时间单元单位毫秒是心跳和超时计算的基础。initLimit是Follower节点启动后与Leader进行初始通信的最大时间10个tickTime也就是20秒。syncLimit是Leader与Follower之间正常通信时的心跳超时时间5个tickTime也就是10秒。关键配置在最后三行server.1/2/3。这表示集群中有三个Zookeeper节点分别对应三台机器。2888端口用于Leader与Follower之间的数据同步3888端口用于Leader选举投票。这两个端口号不能和clientPort冲突也不能和其他应用端口冲突防火墙必须放行。配置完成后在三台机器上分别创建dataDir指定的目录并写入myid文件。这个myid文件里写的就是该节点在配置中的编号比如node01上写1node02上写2node03上写3千万别写错写错了Zookeeper集群起不来。mkdir -p /data/zookeeper echo 1 /data/zookeeper/myid # node01执行 echo 2 /data/zookeeper/myid # node02执行 echo 3 /data/zookeeper/myid # node03执行启动顺序建议先启动node01再node02最后node03。启动方式分别为/opt/zookeeper-3.4.14/bin/zkServer.sh start。启动后通过/opt/zookeeper-3.4.14/bin/zkServer.sh status查看状态。正常情况下三台节点中有1台是Leader另外2台是Follower。如果全是Follower说明选举有问题需要检查网络和myid配置。2.3 三台机器的免密登录配置这一步很多人会跳过但在实际运维中非常重要。Storm的Nimbus在分配任务后需要向Supervisor节点通信并传递代码包如果每次都要输入密码自动化运维和故障恢复就完全没法做。建议在Nimbus所在节点生成密钥对然后把公钥分发到所有Supervisor节点同时也可以配置本机回环免密。ssh-keygen -t rsa ssh-copy-id -i ~/.ssh/id_rsa.pub node01 ssh-copy-id -i ~/.ssh/id_rsa.pub node02 ssh-copy-id -i ~/.ssh/id_rsa.pub node03配置完成后测试ssh node02应该可以无密码登录。3. Storm三大组件安装与核心配置3.1 Storm安装包部署Storm的安装比较轻量本质上就是解压一个压缩包然后修改配置文件。wget https://archive.apache.org/dist/storm/apache-storm-1.2.4/apache-storm-1.2.4.tar.gz tar -zxvf apache-storm-1.2.4.tar.gz -C /opt/ cd /opt/apache-storm-1.2.4/解压后目录结构如下bin/存放启动和管理的脚本conf/存放配置文件lib/存放Storm及依赖的JAR包logs/存放运行日志public/存放Storm UI的前端文件你需要把整个apache-storm-1.2.4目录分发到所有节点保持目录结构一致方便后续管理。我在生产上习惯用rsync同步效率比scp高很多。# 在Nimbus节点执行将storm目录同步到其他节点 rsync -av /opt/apache-storm-1.2.4/ node02:/opt/apache-storm-1.2.4/ rsync -av /opt/apache-storm-1.2.4/ node03:/opt/apache-storm-1.2.4/3.2 storm.yaml配置详解从nimbus.seeds到slot.portsStorm的核心配置全部集中在conf/storm.yaml。这个文件在安装包中已经存在但默认内容很少需要手动补充完整。配置前有个小坑要注意这个文件对缩进极其敏感它解析的是YAML格式空格用错了直接启动失败。网上很多教程给的代码块里缩进乱七八糟照着抄就容易出错。我这里把生产可用的配置贴出来并逐行解释含义。storm.zookeeper.servers: - node01 - node02 - node03 storm.zookeeper.port: 2181 nimbus.seeds: [node01, node02] storm.local.dir: /data/storm supervisor.slots.ports: - 6700 - 6701 - 6702 - 6703 worker.childopts: -Xmx2g -Xms2g ui.port: 8080 storm.messaging.transport: backpressure storm.messaging.netty.buffer.size: 16384 storm.messaging.netty.max.retries: 30 storm.messaging.netty.max.wait.ms: 3000下面逐个说参数含义。storm.zookeeper.servers和storm.zookeeper.port指定Zookeeper集群的地址列表和端口。如果你有多个Zookeeper节点全部列出来即可。nimbus.seeds是Nimbus主节点的候选列表。注意一个非常关键的细节在Storm 1.x及以上版本中这个参数是复数形式nimbus.seeds而网上很多老教程里写作nimbus.host那是0.9.x版本的单Nimbus配置方式。如果你配置的是1.2.x版本写nimbus.host会直接被忽略Nimbus无法启动。这一点最容易坑到人。storm.local.dir是Storm保存本地状态数据的目录。这个目录必须提前创建并且对运行Storm的用户有读写权限。Nimbus会把任务快照存到这里Supervisor会把Worker的元数据存到这里目录权限不对会导致运行时报错。supervisor.slots.ports是关键中的关键。它定义了当前Supervisor节点上可以启动的Worker进程端口列表。举个例子如果配置了4个端口说明这个节点最多同时运行4个Worker进程。每个Worker由一个端口号唯一标识。这个值的设置需要结合机器的CPU核数和内存大小来决定不能一味贪多。worker.childopts是Worker进程的JVM启动参数。-Xmx2g -Xms2g表示每个Worker进程堆内存最大和初始都是2G初始和最大一致的好处是避免JVM在运行期动态扩容导致性能波动。如果你的机器内存充足单Worker堆内存可以调到4G但要注意总内存不能超过物理内存。ui.port是Storm UI的访问端口生产环境中可以改成8088之类或者通过Nginx代理。默认是8080但8080经常被其他服务占用建议提前改掉。下面这段还需要重点解释storm.messaging.transport: backpressure这是一个传输层配置。旧的Storm版本默认依赖ZeroMQ做消息传输但ZeroMQ在部分Linux发行版上的兼容性非常差。从0.10.x开始Storm默认改成Netty作为传输层底层走TCP。但这里要注意storm.messaging.transport这个参数本身在1.2.x版本里有个隐藏的坑如果你配置为netty部分环境会出现消息队列积压问题更稳妥的做法是配置为backpressure这个参数会启用背压机制让Worker在消费不过来时主动减缓上游发送速度避免OOM。3.3 Nimbus节点与Supervisor节点的启动验证配置完成后在三台节点上分别启动对应的服务。Nimbus节点启动NimbusSupervisor节点启动Supervisor还可以在Nimbus节点启动UI服务。# node01、node02作为Nimbus节点 /opt/apache-storm-1.2.4/bin/storm nimbus /dev/null 21 # node01作为UI节点 /opt/apache-storm-1.2.4/bin/storm ui /dev/null 21 # node02、node03作为Supervisor节点 /opt/apache-storm-1.2.4/bin/storm supervisor /dev/null 21 启动后通过jps命令检查Java进程。正常情况能看到一个叫core的进程Nimbus的进程名就是coreUI是core里的一个子模块。如果你执行jps看到QuorumPeerMain这是Zookeeper进程、core说明启动正常。然后用浏览器访问http://node01:8080能看到Storm UI界面。这个界面第一行会显示集群的摘要信息集群中的Supervisor数量、空闲Worker槽位、Topology数量等。如果你看到Supervisor数量是0说明Nimbus还没成功连接上Supervisor节点需要查看Nimbus节点的logs/nimbus.log和Supervisor节点的logs/supervisor.log。注意UI界面显示的是Nimbus视角的集群状态Supervisor节点必须主动向Nimbus注册Nimbus才能感知到它的存在。如果Supervisor启动失败或者网络不通UI界面就一直显示0个Supervisor。3.4 常用Storm命令行操作集群搭好后先测试一下命令是否可用。在Nimbus节点执行/opt/apache-storm-1.2.4/bin/storm list这条命令会和Nimbus通信并列出当前集群中所有运行中的Topology。如果报错Unable to get Nimbus host name说明storm.yaml里的nimbus.seeds配置有问题或者Nimbus进程没有起来。也可以使用storm jar命令来提交任务/opt/apache-storm-1.2.4/bin/storm jar /path/to/my-topology.jar com.example.MyTopology topo-name后面详细讲实际操作。4. 高可用实战与资源参数调优4.1 双Nimbus高可用配置与工作机理在前面配置中我已经写了两个Nimbus节点node01和node02这已经开启了高可用模式。但很多刚接触的人不知道为什么要这么配也不清楚它的工作机制。Nimbus采用的是主备工作模式。两个Nimbus节点会同时启动然后通过Zookeeper进行选主。选出来的主Nimbus负责接收任务提交、分配任务、监控Worker状态等所有调度工作。备用Nimbus节点则处于待命状态持续监控Zookeeper中主节点的心跳信息。如果主Nimbus节点宕机Zookeeper集群会在一个sessionTimeout周期后感知到主节点心跳丢失此时备用Nimbus会被Zookeeper选为新主节点然后接管所有调度任务。这个过程中已经运行中的Topology不会中断因为Worker进程是直接运行在Supervisor节点上的并不依赖Nimbus的持续在线。Nimbus只会在Worker进程故障后重新调度所以宕机时间内只有故障恢复能力短暂缺失不会影响存量任务的数据处理。生产环境里我建议双Nimbus都只部署在独立的控制节点上不要和Supervisor混部。因为Nimbus发生主备切换时会短暂占用大量CPU和内存如果混部在Supervisor节点上可能影响正在运行的Worker进程性能。4.2 CPU与内存的规划Worker槽位数量的计算逻辑这是集群搭建里最容易“拍脑袋”的地方。我先讲一个我设计的容量规划公式再给出实例方便按需套用。假设每台机器有C核CPU内存为MGB操作系统和基础进程预留RGB内存。则规划遵循以下几点每台机器预留的Slot端口数建议等于CPU核数的一半或者等于(M - R) / 每个Worker内存。每个Worker进程建议分配2-4G内存具体取决于你的数据处理逻辑的复杂度和单条消息的大小。举例一台16C32G的服务器运行Supervisor节点操作系统占用4G内存。如果设置每个Worker为2G堆内存可用堆内存是32 - 4 28G那么最多可以开14个Worker。但实际不会开满因为JVM本身还有非堆内存和GC开销所以一般会再打一个折扣建议最多配置10个Worker。supervisor.slots.ports: - 6700 - 6701 - 6702 - 6703 - 6704 - 6705 - 6706 - 6707 - 6708 - 6709对应的worker.childopts设置worker.childopts: -Xmx2g -Xms2g -XX:UseG1GC -XX:PrintGCDetails -Xloggc:/data/gc.log注意如果你的机器内存是32G就开10个Worker。如果你是64G内存也不是无脑加倍建议控制在20个以内。因为Worker进程之间还有网络栈、磁盘IO的竞争进程太多反而增加上下文切换开销吞吐量不一定线性增长。这个我在压测时验证过很多次16C机器上从12个Worker加到16个吞吐几乎没有提升延迟反而略涨。4.3 Topology的并行度规划从源码到集群资源映射很多人搭建完集群后运行Topology发现资源不够报错Not enough supervisor slots for the topology核心原因是没有理解并行度和资源之间的关系。Storm的并行度分为三个层级Worker进程数在代码中使用setNumWorkers(int)指定1个Worker是一个独立的JVM进程。Executor线程数在代码中通过setSpout()和setBolt()的第二个参数指定比如setBolt(bolt1, new MyBolt(), 4)表示这个Bolt有4个并发实例。Task数量不显式设置时Task数等于Executor数一般不需要单独调整。资源映射的逻辑是整个Topology的Worker总数不能超过集群所有Supervisor的空闲Slot总数。而每个Bolt的Executor数可以大于Worker数因为一个Worker进程里可以运行多个Executor线程。举个例子集群有2台Supervisor每台8个Slot共16个Slot。你设置了6个Worker、4个Spout Executor、8个Bolt Executor那么这12个Executor会被调度到6个Worker进程中平均每个Worker跑2个Executor。这里有个最优配置策略尽量让Bolt的并发数等于Bolt运行的Worker数乘以一个合理因子。比如一个Bolt有8个Executor分布在4个Worker上每个Worker跑2个Bolt Executor这就是合理的。如果8个Executor分布在8个不同Worker上每个Worker只跑1个Executor那么一个Worker进程的JVM开销就浪费了资源利用率不高。4.4 背压机制与消息可靠性配置Storm 1.2.x中加入的背压机制Backpressure是个很实用的特性。简单解释就是当某个Bolt处理速度跟不上Spout的发射速度时系统自动反向通知Spout放慢发射速率从而避免Worker内存被撑爆。如果没有背压机制Kafka中积压的消息会通过Spout大量拉取到Worker进程里轻则GC频繁停顿重则Worker直接OOM被杀掉。当初Storm 0.10版本之前经常出这种问题现在配置了storm.messaging.transport: backpressure就不太需要担心了。另外还要说下消息可靠性。Storm的Ack机制默认是开启的每条消息被处理完后都会通过AckTracker来确认。但这会带来额外的性能开销。如果你做的是实时指标计算允许消息偶尔丢失可以关闭Ack机制把setNumAckerExecutors()设置为0。但如果你做的是金融交易类的实时风控每条消息都不能丢那就必须开启Ack而且还要配合Kafka的offset手动管理来实现精确一次语义。这个根据业务取舍。5. 提交测试Topology与集群状态验证5.1 用自带Example做端到端验证集群搭好之后第一步跑一个官方自带的示例程序做全链路验证。Storm安装包里自带了一个storm-starter的例子但如果你是手动下载的压缩包里面不一定有完整的示例JAR。更稳妥的方法是自己写一个极简Topology提交上去。下面是一个最简单的Spout-Bolt示例它每秒发射一条消息然后打印接收到的内容。import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.StormSubmitter; import org.apache.storm.topology.TopologyBuilder; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.topology.base.BaseRichSpout; import org.apache.storm.task.OutputCollector; import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import java.util.Map; import java.util.UUID; import java.util.concurrent.atomic.AtomicLong; public class TestTopology { public static class TestSpout extends BaseRichSpout { private OutputCollector collector; private AtomicLong counter new AtomicLong(0); public void open(Map conf, TopologyContext context, OutputCollector collector) { this.collector collector; } public void nextTuple() { String msg message- counter.incrementAndGet(); collector.emit(new Values(msg)); try { Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); } } public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(message)); } } public static class PrintBolt extends BaseRichBolt { private OutputCollector collector; public void prepare(Map conf, TopologyContext context, OutputCollector collector) { this.collector collector; } public void execute(Tuple tuple) { String message tuple.getStringByField(message); System.out.println(Received: message); collector.ack(tuple); } public void declareOutputFields(OutputFieldsDeclarer declarer) { } } public static void main(String[] args) throws Exception { TopologyBuilder builder new TopologyBuilder(); builder.setSpout(test-spout, new TestSpout(), 1); builder.setBolt(print-bolt, new PrintBolt(), 1) .shuffleGrouping(test-spout); Config config new Config(); config.setNumWorkers(2); if (args ! null args.length 0) { StormSubmitter.submitTopology(args[0], config, builder.createTopology()); } else { config.setMaxTaskParallelism(1); LocalCluster cluster new LocalCluster(); cluster.submitTopology(test, config, builder.createTopology()); Thread.sleep(30000); cluster.shutdown(); } } }编译打包后提交到集群mvn clean package -DskipTests /opt/apache-storm-1.2.4/bin/storm jar target/test-topology-1.0.jar TestTopology test-topology提交后通过Storm UI查看Topology状态。在UI界面左侧能看到Topology列表点进去可以看到Spout和Bolt的并发数、延迟、发送的消息数等指标。如果一切正常在Supervisor节点的logs/worker-*.log中能看到Received: message-1、Received: message-2这样的输出。5.2 读懂Storm UI的关键指标Storm UI是判断集群健康度的第一入口。我在这里分享一下我每天查看UI时必看的几个指标按优先级排列指标正常参考值说明Completed持续增加表示Bolt成功处理的消息数如果停止增长说明有地方挂住了Failed0或极少超过5%就要排查可能Spout或者Bolt内部抛异常Capacity低于1.0表示Executor的处理能力是否饱和大于1说明需要扩容Executor延迟视业务而定如果持续超过几秒说明Bolt里访问外部存储的耗时太长Worker心跳无超时如果出现Timed out表示Worker被OOM或者卡死特别注意Capacity这个指标。它的含义是该Executor处理消息的总时间除以运行时间的比值。比如一个Executor运行了100秒其中80秒在处理消息20秒在空闲Capacity就是0.8。如果Capacity长期大于1说明这个Executor已经积压大量消息处理不过来你需要提高这个Bolt的并发数或者优化Bolt内部的逻辑。5.3 日志查看Nimbus/Supervisor/Worker三层日志的定位方法排查Storm问题最依赖的就是日志。日志分布在两个目录logs/和logs/workers-artifacts/。前者存在Nimbus和Supervisor进程的日志后者存在每个Topology的每个Worker的详细日志。logs/nimbus.logNimbus的运行日志里面包含任务调度的记录、心跳检查、异常信息。logs/supervisor.logSupervisor的运行日志可以看到Worker的启动和关闭记录。logs/workers-artifacts/topology-id/worker-port/worker.log具体Worker进程日志里面是Topology中Spout/Bolt的System.out.println输出和异常堆栈。查看Worker日志的方法tail -f /opt/apache-storm-1.2.4/logs/workers-artifacts/test-topology-*/6700/worker.log这个文件是排查Topology运行时异常的第一现场任何Spout或Bolt里的未捕获异常都会打印到这里。6. 生产环境扩展整合Kafka与HBase6.1 数据实时链路设计与组件联动一个完整的大数据实时处理平台不会只有Storm一个组件。我生产环境里最常见的实时链路典型架构如下Kafka - Storm (Spout消费) - Bolt处理 - HBase / Redis / 其他存储Kafka作为数据缓冲层负责承接各种上游业务系统的数据。Storm的KafkaSpout负责从Kafka的指定Topic中拉取数据然后发射到Bolt中进行业务计算最后写入HBase或者Redis。在这个架构里Kafka的Topic分区数决定了Storm的最大并行度上限。一个原则是Kafka的Topic分区数应该大于或等于Storm中Spout的并发数。否则会出现多个Spout线程争抢同一个Kafka分区造成数据重复消费或分区间负载不均。比如你的Topic有12个分区Spout并发可以设置为12或6通常建议等于分区数除非上游数据量很小。6.2 KafkaSpout配置与Offset管理要点KafkaSpout的配置是实时链路里最需要精心调优的地方。下面给出一个在Storm 1.2.x中集成Kafka 2.x的配置片段。import org.apache.storm.kafka.spout.KafkaSpout; import org.apache.storm.kafka.spout.KafkaSpoutConfig; KafkaSpoutConfigString, String kafkaConfig KafkaSpoutConfig.builder(node01:9092,node02:9092,node03:9092, input-topic) .setProp(group.id, storm-group) .setProp(max.poll.records, 500) .setProp(enable.auto.commit, false) .setProp(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer) .setProp(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer) .setFirstPollOffsetStrategy(KafkaSpoutConfig.FirstPollOffsetStrategy.UNCOMMITTED_EARLIEST) .setProcessingGuarantee(KafkaSpoutConfig.ProcessingGuarantee.AT_LEAST_ONCE) .build(); KafkaSpoutString, String kafkaSpout new KafkaSpout(kafkaConfig);这里有几个核心参数需要注意max.poll.records决定一次Kafka拉取的消息条数。如果单条消息很大建议调小如果消息都很小可以适当调大减少网络轮询开销。enable.auto.commit必须设置为false让Storm的Ack机制来控制offset提交保证消息不丢不重由Storm在确认消息被完整处理后提交offset。setProcessingGuarantee选择AT_LEAST_ONCE即可保证不丢但可能重复。如果要做精确一次需要配合Kafka的幂等Producer和事务机制但复杂度高一个量级。setFirstPollOffsetStrategy设成UNCOMMITTED_EARLIEST意思是从未提交的offset开始消费如果没有历史offset则从最早的开始这个策略适合新上线链路。6.3 Bolt中批量写入HBase的优化方案实时计算最容易出现的性能瓶颈不在计算而在结果写入外部存储。如果每条消息都单独写一次HBaseHBase的RegionServer压力会非常大而且会因为网络往返损耗导致吞吐量上不去。常见的优化方案是在Bolt中维护一个批量写入缓冲区累积一定条数或者一定时间后批量提交。伪代码如下public class HBaseWriteBolt extends BaseRichBolt { private ListPut buffer new ArrayList(); private OutputCollector collector; private Connection conn; private static final int BATCH_SIZE 100; private static final long BATCH_TIMEOUT_MS 2000; private long lastFlushTime; public void execute(Tuple tuple) { Put put buildPut(tuple); buffer.add(put); collector.ack(tuple); if (buffer.size() BATCH_SIZE || System.currentTimeMillis() - lastFlushTime BATCH_TIMEOUT_MS) { flushToHBase(); } } private void flushToHBase() { // 批量put到HBase使用HTable的batch接口 // table.batch(buffer) buffer.clear(); lastFlushTime System.currentTimeMillis(); } }这个思路需要特别注意批量写入虽然提高了吞吐但如果在批量提交之前Worker宕机这批数据就会丢失。所以在生产环境中批量提交方案必须配合Kafka的offset提交方式来权衡。如果精确一次性语义是这个业务的硬性需求就不能用这种内存批量缓冲需要考虑在Bolt内用事务或幂等写入来规避。7. 集群运行常见问题与故障排查速查表7.1 Nimbus启动失败排查症状执行storm nimbus后进程秒挂或者storm list报错无法连接。排查步骤查看logs/nimbus.log重点看最后50行。如果出现Cannot connect to Zookeeper说明Zookeeper集群有问题先检查Zookeeper是否正常三台节点之间是否网络通。如果出现YAML syntax error说明storm.yaml中缩进或格式错误逐行检查配置。确认nimbus.seeds配置的是IP还是主机名。如果配置主机名必须在/etc/hosts中配置所有节点的IP映射。这个场景最常见的两个原因一个是nimbus.host和nimbus.seeds混用另一个是storm.local.dir目录权限不够。在1.2.x版本中写nimbus.host完全没用必须用nimbus.seeds。7.2 Supervisor注册失败排查症状Nimbus UI界面Supervisor数量为0或Supervisor状态显示Disconnected。排查思路在Supervisor节点上确认storm supervisor进程已启动。查看logs/supervisor.log如果出现Old timestamps found说明该节点的本地状态和Nimbus不一致可以删除storm.local.dir下的supervisor目录然后重启Supervisor。检查Supervisor节点到Nimbus节点及Zookeeper节点的网络连通性开放8080、6627、6700-6710等端口。如果有防火墙要确保net.python等通信端口都放行。这里有个值得注意的细节如果删除storm.local.dir下的目录该节点上所有已经分配的任务会被清空Nimbus会重新调度这些任务到其他可用节点。所以如果是生产环境不要在业务高峰期做这个操作。7.3 Worker进程持续重启的排查症状Topology的某个Bolt的Executor迟延率高Worker频繁重启UI界面显示Worker心跳超时。排查思路查看logs/workers-artifacts/topology-id/port/worker.log看是否有OutOfMemoryError。如果有OOM把worker.childopts中的-Xmx2g调大到-Xmx4g或者减小该节点的Slot数减轻单个进程压力。查看是否有外部依赖连接超时比如HBase连接池耗尽、Redis连接不可用。这类异常通常表现为Connection refused或SocketTimeoutException。确认Bolt代码中没有死循环或无限等待锁等逻辑。这类问题排查难度较高建议在代码中加入日志和超时控制。7.4 消息积压与处理延迟升高症状Kafka的消费组Lag持续增加Storm UI中Bolt的延迟时间不断增长。排查思路先看是单Bolt瓶颈还是全链路瓶颈。在Storm UI中分别看每个Bolt的延迟和Capacity指标。如果只有某个Bolt延迟高说明该Bolt的处理逻辑有问题或者外部依赖如数据库响应慢需要优化Bolt代码或者增加该Bolt的并发数。如果是全链路都慢说明集群整体资源不足需要增加Supervisor节点或者增加Worker数量同时检查CPU和内存是否被打满。看看是否出现反压。如果storm.messaging.transport配置为backpressureSpout会主动放慢拉取速度此时Kafka Lag升高但Worker本身不会OOM这其实是系统在自我保护需要做的是扩容。7.5 常用故障排查命令速查表这里整理了生产环境里最常用的一批命令贴在下面方便查阅。场景命令查看进程状态jps -l查看Nimbus日志tail -f /opt/apache-storm-1.2.4/logs/nimbus.log查看Supervisor日志tail -f /opt/apache-storm-1.2.4/logs/supervisor.log查看Worker日志tail -f /opt/apache-storm-1.2.4/logs/workers-artifacts/topology/port/worker.log查看端口占用lsof -i:6700查看Zookeeper状态/opt/zookeeper-3.4.14/bin/zkServer.sh status查看集群Topology列表/opt/apache-storm-1.2.4/bin/storm list动态调整Topology并行度在Storm UI上点击Topology的rebalance按钮8. 基于实际维护经验的几点建议走完整个搭建流程最后分享几个我踩过的坑和长期沉淀下来的维护理念。第一配置文件必须版本管理。storm.yaml、zoo.cfg、JDK安装路径这些都属于基础设施应该纳入Git仓库统一管理并配合Ansible或Shell脚本一键部署。不然半年后机器故障需要重新搭建你根本想不起来当时改了哪些参数。第二监控必须早于故障。我个人建议集群上线之前就要对接监控告警体系。至少要监控以下指标Zookeeper的存活数和Leader切换次数Nimbus的存活状态最好做成自动拉起机制每个Supervisor节点上Worker的存活数每个Topology的Completed、Failed、Capacity指标Kafka消费组的Lag值一旦Topology的Failed数据超过阈值或者Kafka Lag突然飙升立即触发告警而不是等业务方反馈说数据延迟了。第三进程启动脚本要加日志重定向。很多教程里nohup启动Storm进程时不重定向输出导致启动后看不到任何输出。但在生产环境这个习惯很危险。建议用如下方式启动并保留控制台输出方便排查启动阶段的问题nohup /opt/apache-storm-1.2.4/bin/storm nimbus /opt/apache-storm-1.2.4/logs/nimbus-console.log 21 第四升级版本前先在测试环境压测。Storm本身是Java进程依赖的Netty、Zookeeper客户端库版本都可能影响运行稳定性。我见过从Storm 1.1升到1.2后因为没有同步升级KafkaClient插件导致新旧序列化协议不一致线上数据大面积反序列化失败的案例。版本升级不是小事每一步都要有明确的兼容性验证。Storm的集群搭建本身不难难的是理解它背后的协调逻辑并能在出现故障时快速定位问题。希望这篇实操笔记能帮你少走几个弯路搭出一个健壮的集群后续把实时计算业务稳定跑起来。