
简介基于Java与Apache Storm的日志监控告警系统面向后端、大数据及运维监控方向开发者解决实时消费Kafka日志、按规则识别异常并触发邮件/短信告警的需求。项目以Storm拓扑串联数据接入、规则加载、日志处理、通知分发和告警入库等核心环节并封装CommonUtils、JdbcUtils等工具类便于理清实时流处理与外部存储、消息队列的协作方式。压缩包共100个文件其中24个Java源码为拓扑各Bolt具体实现27个Class为编译产物另有XML配置文件、MD说明文档和39张PNG流程截图包体仅1.17MB轻量完整可快速对照源码和图示梳理运行过程。已有159人学习浏览适合希望通过完整样例掌握Storm消费Kafka、规则匹配与多渠道告警的开发者也可作为扩展规则引擎和告警模块的参考基础。1. 日志监控告警系统.zip一个能直接改的StormKafka实时告警骨架做运维监控的人常常会遇到一种尴尬日志量不大时用ELK硬扛也能出告警业务一上来Kafka里堆了几百万条日志告警延迟从秒级变成分钟级。我拆的这个“基于JavaStorm的日志监控告警系统.zip”核心就是解决这个问题——用Apache Storm拓扑从Kafka实时消费日志经过规则匹配后把异常通过邮件和短信推出去同时落库记录。压缩包里的类名很完整TopologyMain、StormTickBolt、TopkeyCountBolt、JdbcUtils、ShortMessageUtil……从这些类名能看出它不是一个只写概念的Demo而是一个把消费、计算、通知、存储串起来的生产骨架。适合两类人正在做实时告警平台、想抄一条完整链路的人以及想搞懂Storm和Kafka怎么配合的Java工程师。如果你只想要概念这篇文章会讲清楚每个Bolt的作用和参数如果你要把它跑通后面有完整步骤和踩坑记录。2. 拆解Storm拓扑从Kafka Spout到告警落地的数据流日志监控系统听上去复杂但真正跑起来就是一条流水线Kafka负责堆日志Storm负责搬货和分拣分拣发现异常就发邮件发短信同时记一笔账。这套架构的好处是Storm里的每个Bolt只干一件小事出问题可以单独扩容不用把整个消费逻辑闷在一个线程池里。2.1 从Kafka Spout到告警发送这条流水线上的七个关键类拿这个项目来说数据流大致是Kafka Spout → ProcessDataBolt规则匹配 → NotifyMessageBolt邮件/短信 / SaveToDBBolt落库另外还有一个StormTickBolt专门负责定时刷新规则。压缩包里没有把所有源码列出来但通过类名和摘要描述每个节点的职责非常清楚。类名在项目里的实际角色对应数据流位置TopologyMain主拓扑入口组装Spout与Bolt起点TopkeyTopologMainTopKey统计拓扑入口与TopkeyCountBolt配套可选第二个拓扑StormTickBolt按固定频率从MySQL加载监控规则、应用和用户信息规则广播源ProcessDataBolt解析日志与预加载规则做匹配核心处理节点NotifyMessageBolt根据匹配结果发送邮件/短信通知分支SaveToDBBolt把告警记录写入数据库存储分支CommonUtils规则匹配、通知发送等公共方法被多个Bolt调用JdbcUtils数据库连接和数据读取供StormTickBolt等使用LogMonitorUser用户信息实体给NotifyMessageBolt提供收件人注意项目正文里直接给出的.class文件里没有ProcessDataBolt和NotifyMessageBolt但摘要里描述了这两个Bolt它们大概率放在别的包中。其它类像MessageSenderUtil、ShortMessageUtil、MailInfo明显是通知模块的一部分。TopkeyCountBolt和TopkeyTopologMain则说明这套系统不止一个拓扑主拓扑做实时告警另一个拓扑可能按业务Key统计告警频率或日志量。2.2 为什么选Storm流式计算与Kafka消费的边界有人会问直接用Kafka Consumer加线程池不行吗当然行但你会发现要自己管理的事情变多消费线程挂了怎么办消息处理失败要不要重试同一批规则如何共享到所有线程窗口统计怎么做。这些东西自己做不是不能但做出来的东西很难达到Storm这种分布式流处理框架的成熟度。Storm在这里的角色是把Kafka Consumer包成一个KafkaSpout消息进来后交给后续Bolt。Spout只负责“读”Bolt只负责“算”。比如ProcessDataBolt拿到一条日志用CommonUtils里的规则集合做关键字匹配匹配到了就把告警信息emit给下游。这个模型在Java面试题里也常被问到Kafka和Storm、ZooKeeper如何协作Storm的ack/fail机制怎么保证不丢数据。选Storm而不是Flink对这个项目来说倒不是技术优劣问题。这个包明显是一个以Storm为中心的工程而且Storm的模型对“日志进来→判断→通知”这种轻量计算非常契合延迟可以做到毫秒级。Flink的优势体现在复杂窗口、状态管理和Exactly-Once语义但如果只是规则匹配加发通知Storm的At-Least-Once加上外部幂等就能满足大部分场景。2.3 Topology装配代码怎么写setSpout与setBolt的参数实际工程里拓扑装配集中在TopologyMain。如果你手里只有编译后的.class用JD-GUI反编译后看到的逻辑也八九不离十。我这里写一个简化版对应这个项目的主链路// TopologyMain 核心代码KafkaSpout - ProcessDataBolt - NotifyMessageBolt / SaveToDBBolt TopologyBuilder builder new TopologyBuilder(); // Kafka Spout从 app_log 主题读日志 KafkaSpoutConfigString, String spoutConfig KafkaSpoutConfig .builder(localhost:9092, app_log) .setGroupId(log-monitor-group) .setFirstPollOffsetStrategy(EARLIEST) .build(); // 并行度设为2大致匹配Kafka分区数 builder.setSpout(kafka-spout, new KafkaSpout(spoutConfig), 2); // ProcessDataBolt并行度设为4按level字段分组保证同一级别日志进同一个task builder.setBolt(process, new ProcessDataBolt(), 4) .fieldsGrouping(kafka-spout, new Fields(level)); // 通知Bolt随机接收process发来的告警元组 builder.setBolt(notify, new NotifyMessageBolt(), 2) .shuffleGrouping(process); // 存储Bolt同样随机接收 builder.setBolt(save, new SaveToDBBolt(), 2) .shuffleGrouping(process); Config config new Config(); config.setNumWorkers(3); config.setDebug(false); // 让 tick 每 60 秒触发一次用于规则刷新 config.put(Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS, 60);这段代码里有几个参数需要你真正去调KafkaSpoutConfig.builder()的第一个参数是broker地址列表多个broker用逗号分隔第二个参数是要消费的topic名。setFirstPollOffsetStrategy(EARLIEST)决定拓扑启动时从最早offset开始读还是从最新开始如果规则匹配逻辑依赖历史数据用EARLIEST如果只关心启动后的日志用LATEST。fieldsGrouping(kafka-spout, new Fields(level))表示按日志级别INFO/ERROR把消息路由到同一个处理task这样后续做计数统计时不会乱。setNumWorkers(3)是启动3个Worker进程不是3个JVM线程调试时建议改成1不然日志分散在多个worker里很难看。规则刷新Bolt我一般不让它挂在主线数据流上而是通过一个静态Map共享规则快照。StormTickBolt固定1个并行度靠Config里的tick配置每60秒醒来一次从数据库加载规则写入Map。ProcessDataBolt执行时直接读Map效率远高于每条日志都查一次数据库。实际项目里有人非要把规则刷新Bolt用allGrouping接到主链路上结果每次日志进来都会触发一次广播反而给网络造成不必要的负担。2.4 背压与消息确认At-Least-Once下为什么不能随便ack这个项目的告警通知走的是邮件和短信而短信接口一般都不是百分百可靠。你在Bolt里调用MessageSenderUtil发送短信时如果发送超时最好抛异常并让框架重发而不是吞掉异常后ack。否则这条告警就永远丢了。但反过来说At-Least-Once带来的副作用是重复消息。短信发送成功后Spout可能因为网络抖动没有收到ack会把同一条日志重新发下来。所以告警通知Bolt里必须有去重同一个ruleId加同一个日志指纹在时间窗口内只发一次。后面避坑部分我会专门讲这条。3. 把项目跑起来从解压到Kafka联调的实操步骤拿到这个zip你的目标不是看一遍类名而是要让它在你本地或服务器上转起来。按下面的流程走每一步都有对应的验证方式。3.1 解压后先看目录结构别急着导入IDE先解压确认这个包是源码工程还是纯编译产物unzip 基于JavaStorm的日志监控告警系统.zip -d log-monitor cd log-monitor find . -name *.class | wc -l find . -name pom.xml -o -name .classpath | head -5这段命令先解压到log-monitor目录然后统计.class文件数量再判断是Maven工程还是Eclipse工程。如果.class数量多但没有pom.xml说明压缩包里是编译好的产物需要用JD-GUI或FernFlower反编译成Java源码。如果连class文件都没有只有Java文件那恭喜可以直接走正常导入流程。我见过有同事拿到zip直接拖进IDE结果一堆Cannot resolve symbol原因就是没先确认是源码还是class包。反编译这种事说穿了也就两分钟JD-GUI里打开jar或单个.classCtrlA导出全部源码。反编译出来的代码可能变量名变成var1但逻辑骨架还在。这个项目里CommonUtils和JdbcUtils是关键先反编译这两个类比乱猜有用得多。3.2 本地环境搭配不折腾版本的人会少踩一半坑日志监控告警这种系统版本搭配非常重要。我试过用Storm 2.x配Kafka 2.8结果KafkaSpout的依赖一直冲突最后还是换回老组合才跑通。这个项目是基于Java和Apache Storm的常见的稳定搭配如下组件推荐版本区间说明JDK1.8Storm 1.x系列在JDK8下最稳高版本要小心反射权限ZooKeeper3.4.xStorm1.x用ZK做协调版本别太新Kafka2.2 ~ 2.5与Storm KafkaClient兼容性较好Storm1.2.x1.x版本资料最多踩坑容易搜到解决方案设置好环境变量再跑工程export JAVA_HOME/usr/lib/jvm/java-8 export STORM_HOME/opt/apache-storm-1.2.2 export PATH$PATH:$STORM_HOME/bin说明一下Kafka虽然名叫Kafka但和Storm集成时还需要引入storm-kafka-client依赖版本要和Storm主版本一致。如果你用Maven就在pom.xml里加对应依赖如果反编译出来是纯class可以用java -cp手动引入所有依赖jar但那样太痛苦建议还是用Maven重建工程。3.3 数据库初始化建规则表并造一条监控规则这个系统离不开数据库。先按下面SQL建两张表一张存监控规则一张存告警记录CREATE TABLE monitor_rule ( rule_id INT PRIMARY KEY AUTO_INCREMENT, rule_name VARCHAR(64) NOT NULL, keyword VARCHAR(128) NOT NULL COMMENT 规则匹配关键字, severity TINYINT DEFAULT 1 COMMENT 告警级别 1-3, notify_type VARCHAR(32) DEFAULT email, app_id INT DEFAULT 0 ); CREATE TABLE alert_record ( id BIGINT AUTO_INCREMENT PRIMARY KEY, app_name VARCHAR(64), rule_name VARCHAR(64), content TEXT, alert_time DATETIME, KEY idx_alert_time (alert_time) );规则表里的keyword是核心比如你想监控“timeout”和“OutOfMemoryError”就往monitor_rule里插两条记录。JdbcUtils一定是从库里查这些规则然后给StormTickBolt加载进内存。启动系统前先插入测试规则INSERT INTO monitor_rule (rule_name, keyword, severity, notify_type) VALUES (超时告警, timeout, 2, email,short_message); INSERT INTO monitor_rule (rule_name, keyword, severity, notify_type) VALUES (OOM告警, OutOfMemoryError, 3, email,short_message);注意alert_time字段一定要建索引。因为后续你会经常查“最近10分钟有哪些告警”没有索引的话落库数据一多这条查询会直接把数据库拖死。3.4 启动Kafka并模拟日志流验证全链路确认数据库OK后先启动Kafka再创建topickafka-topics.sh --create \ --topic app_log \ --partitions 3 \ --replication-factor 1 \ --bootstrap-server localhost:9092partitions设置3是因为我们拓扑里Spout并行度也是2到3如果分区数远大于Spout并行度有些线程会空闲如果分区数小于并行度又会有线程空转。然后启动一个console producer往里塞一条模拟日志kafka-console-producer.sh \ --topic app_log \ --bootstrap-server localhost:9092输入下面这条JSON格式日志后回车{app:order-service,level:ERROR,msg:[worker-1] connection timeout when calling pay-service}这条日志包含了关键词ERROR和timeoutProcessDataBolt匹配后会emit到一个告警对象NotifyMessageBolt就会尝试发邮件和短信。如果你没有配邮件服务器就把notify_type改成只存库先验证落库那一段。3.5 提交拓扑到Storm集群并观察worker日志本地调试可以用LocalCluster替身但正式跑还是得提交到Storm集群执行storm jar log-monitor.jar com.example.TopologyMain log-monitor-topology storm list服务端提交后用storm list看拓扑状态是否是ACTIVE。如果状态是ACTIVATE而不是ACTIVE说明拓扑还没起来通常是Kafka连不上或ZooKeeper会话超时。然后跟踪worker日志查看异常storm logs日志里有几条常见输出KafkaSpout拉取消息的debug信息ProcessDataBolt匹配成功的info日志以及NotifyMessageBolt发送失败时的exception stacktrace。我一般不看完整堆栈只看Caused by部分十有八九是数据库连不上、短信接口超时、或者反序列化失败。4. 避坑规则匹配与告警策略里最容易翻车的六个问题规则匹配看起来简单就是if contains else continue但把规则放进Storm拓扑后就容易出现各种奇怪问题。这一章我整理了五条真实踩坑记录每一条我都见过不止一次。4.1 规则匹配的常见设计关键字命中与正则的取舍这个项目里规则大概率存在数据库里每个规则有一个keyword字段。最简单的匹配方法是log.contains(rule.getKeyword())。但如果你面对的是复杂日志比如要匹配“IP访问频率超过100次”单纯contains就不够用了。我一般会在CommonUtils里保留一个matchRule(String content, MonitorRule rule)方法里面先用contains做粗筛再对需要正则的规则做Pattern.matches()。正则规则要单独存字段不要和keyword混在一起。因为正则编译非常耗CPUBolt里每条消息都动态compile会直接把Worker打满。// CommonUtils 里典型的规则匹配方法示意 public static boolean matchRule(String content, MonitorRule rule) { if (rule.getKeyword() ! null content.contains(rule.getKeyword())) { return true; } if (rule.getRegexp() ! null rule.getRegexp().length() 0) { // 正则要提前编译好缓存到Map里避免每次execute都compile Pattern pattern patternCache.get(rule.getRuleId()); if (pattern null) { pattern Pattern.compile(rule.getRegexp()); patternCache.put(rule.getRuleId(), pattern); } return pattern.matcher(content).find(); } return false; }这段代码的逻辑说明先做关键字包含匹配命中直接返回true如果规则里有正则就从缓存里拿编译好的Pattern避免反复编译。参数content是Kafka消费出来的日志原文rule是从数据库加载的规则对象。如果你要增加规则只需要往monitor_rule表里插入新行不用改代码。4.2 StormTickBolt的定时刷新别让每个Bolt都查数据库这个系统里StormTickBolt的存在意义就是解决“规则什么时候更新、怎么同步到每个task”的问题。如果每个ProcessDataBolt在execute里查一次数据库那Kafka来一条消息就查一次Kafka堆积几万条时数据库直接崩。正确做法是让StormTickBolt以固定频率查库把规则集合更新到本地内存再通过ZooKeeper或广播流发给其他Bolt。更简单的做法是把它写成一个静态MapStormTickBolt只负责putProcessDataBolt只负责get。启动时把Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS设为30或60表示每30秒触发一次TickTuple。4.3 五条踩坑记录现象、原因、解决方案第一条Kafka一直有数据但拓扑里Bolt收不到消息。现象是storm list显示拓扑运行正常但数据库里一条告警都没新增。原因是KafkaSpout的group.id和另一个消费组冲突或者offset被提交到了无法访问的路径。解决方法是检查KafkaSpoutConfig中的setGroupId保证唯一然后用kafka-consumer-groups.sh --describe --group log-monitor-group查看lag。如果lag持续为0但数据没进Bolt多半是反序列化失败消息被静默drop了。第二条短信被重复轰炸一个故障发出几十条。原因是NotifyMessageBolt在调用短信接口时线程阻塞超时导致Bolt返回failKafkaSpout重新发送同一条日志。解决方法是两个第一条是发送成功后立即ack并记录告警指纹第二是增加发送频率控制比如每个ruleId每分钟最多发一次。短信接口响应慢时别用同步调用把消息转发到一个内部异步队列会好很多。第三条修改规则后迟迟不生效必须重启拓扑。原因是StormTickBolt虽然定时拉库但ProcessDataBolt里的静态Map没有被更新。解决方法是确认Bolt里拿的是同一个Map引用而不是每次执行都new一个Map。另外检查tick配置是否真的生效——在本地调试时LocalCluster不会自动加载tick配置必须手动config.put。第四条Topology启动后频繁报连接池耗尽。现象是worker日志大量出现Connection is not available, request timed out。原因是JdbcUtils底层使用DriverManager.getConnection()每次execute都新建连接高并发下连接数爆炸。解决方法是换成HikariCP或Druid连接池在Bolt的prepare()里只初始化一次execute里直接复用。第五条计数结果乱套TopkeyCountBolt统计的key数量不对。现象是同样的关键字一分钟前统计是100一分钟后统计是80数据反而倒退了。原因是fieldsGrouping分组的字段大小写不一致比如日志里一会儿是ERROR一会儿是error导致相同日志被路由到不同task。解决方法是在Spout或ProcessDataBolt入口统一字段值例如一律转成大写level.toUpperCase().trim()。5. 数据存储与二次分析告警记录不只是看一遍很多项目把告警存到数据库就完事了实际上这些记录是大宝藏。通过分析告警频率你能找到系统真正的弱点在哪。所以落库这步不是随便insert一下架构上要留出查询余地。5.1 落库表结构告警记录表要怎么设计才不后悔前面给出的alert_record表是个基础版本生产环境我还会加几个字段status待处理/已处理/忽略、source_app哪个应用来的、alert_hash去重指纹。alert_hash特别重要它的值可以是ruleId 日志摘要的MD5在NotifyMessageBolt发送前先用这个hash查redist或数据库如果5分钟内已经存在就直接跳过。5.2 JdbcUtils与连接池把DriverManager换成HikariCP的完整替换这个项目原生的JdbcUtils如果是用DriverManager在低并发演示没问题一旦日志量大就撑不住。换成HikariCP的改动其实非常小保留JdbcUtils的对外方法内部实现替换成连接池// JdbcUtils 连接池改造示意 HikariConfig config new HikariConfig(); config.setJdbcUrl(jdbc:mysql://localhost:3306/log_monitor); config.setUsername(root); config.setPassword(123456); config.setMaximumPoolSize(20); config.setMinimumIdle(2); config.setConnectionTimeout(3000); HikariDataSource dataSource new HikariDataSource(config); // 原getConnection方法内部改为 dataSource.getConnection()这里说下参数匹配maximumPoolSize不要设太大Storm一个worker通常会有多个Bolt线程每个Bolt持有自己的连接引用如果整个拓扑并行度是42220个连接足够了。connectionTimeout设3000毫秒超过3秒直接报错而不是无限等。我见过有人把连接超时设为60秒结果数据库一慢大量线程卡在getConnection上拓扑假死。5.3 用TopkeyCountBolt做告警频率聚合从原始日志到TopNTopkeyCountBolt这个类从名字看是在做TopKey统计。它可以统计某一类异常在时间窗口内出现的次数超过阈值就升级告警级别。这个Bolt的实现逻辑很简单维护一个HashMap计数然后在tick tuple触发时把TopN记录往下游发// TopkeyCountBolt 计数逻辑示意 private MapString, Integer countMap new HashMap(); public void execute(Tuple tuple) { if (isTickTuple(tuple)) { // 输出当前窗口内次数最多的Top10异常key for (Map.EntryString, Integer e : topN(10)) { collector.emit(new Values(e.getKey(), e.getValue())); } countMap.clear(); } else { String key tuple.getStringByField(ruleName); countMap.merge(key, 1, Integer::sum); } }这个Bolt的价值在于把告警从“单条命中”升级为“高频聚合”。比如某个规则每分钟只允许触发一次如果一秒钟内被触发100次说明系统正处于一次大规模故障中这时应该发一个P0级短信而不是每条都发低级通知。执行这个Bolt时要注意tick tuple和普通tuple要分开处理只用isTickTuple判断tuple来源千万不要把tick也当普通数据塞进countMap。5.4 告警查询SQL十分钟内到底发生了什么最后再给一个实用查询每次系统出故障我都会先跑这条SQL看最近10分钟的告警趋势SELECT rule_name, COUNT(*) AS cnt, MIN(alert_time) AS first_time, MAX(alert_time) AS last_time FROM alert_record WHERE alert_time NOW() - INTERVAL 10 MINUTE GROUP BY rule_name ORDER BY cnt DESC这条SQL能快速看出哪个规则在爆发以及故障持续了多久。如果你发现cnt特别大就去对应应用日志里找根因如果cnt很小但系统也异常说明你的监控规则没覆盖到真正的故障点这时候需要回头调整规则而不是改代码。6. 进阶生产级改造前要自查的五个习惯这个项目能跑通只是第一步要拿到生产环境用我建议你每次提交拓扑前都检查五个点。第一个是拓扑的并行度与Kafka分区数是否匹配。Spout的并发不能远大于分区数否则会有线程空转也不能远小于分区数否则Kafka的吞吐优势发挥不出来。第二个是规则数据库的连接池是否够用尤其是StormTickBolt每30秒查一次库如果连接池太小tick刷新就会挤压业务查询。第三个是告警通知的幂等机制短信和邮件一定要引入去重表否则一次故障就能让你收到几十条重复短信。第四个是日志和告警数据的时间戳统一用服务器时间不要用日志自带时间因为应用服务器和监控服务器的时间很容易差出几秒导致窗口统计不准。第五个是每个Bolt的失败重试要区分业务错误和环境错误短信接口返回余额不足重试一万次也没用网络超时则应该立刻fail并进入重发队列。我自己的习惯是每次写完Storm拓扑先启动一个模拟日志生产者灌十条包含异常关键字的日志到Kafka然后盯着kafka-consumer-groups.sh里的lag等lag清零后再看数据库alert_record表里是否新增了对应记录。这条链路完整走通后才会把拓扑提交到集群生产环境。否则你根本不知道是Kafka的问题、规则的问题还是拓扑装配的问题一切都在黑匣子里排查起来极其痛苦。有一次我就是没做这一步结果上线后才发现NotifyMessageBolt的短信接口配置被写死了测试账号整个下午所有告警都发到一个不存在的号码上数据库里却全是记录看起来一切正常。从那以后我每次都强制走一遍“模拟日志→检查offset→观察worker日志→验证告警落地”这四步希望帮到你。本文还有配套的精品资源点击获取