☰
基于Storm+Kafka的日志实时监控告警系统架构与实战
2026/10/7 13:23:41 网站建设 项目流程

简介:基于Java与Apache Storm的日志监控告警系统,面向后端、大数据及运维监控方向开发者,解决实时消费Kafka日志、按规则识别异常并触发邮件/短信告警的需求。项目以Storm拓扑串联数据接入、规则加载、日志处理、通知分发和告警入库等核心环节,并封装CommonUtils、JdbcUtils等工具类,便于理清实时流处理与外部存储、消息队列的协作方式。压缩包共100个文件,其中24个Java源码为拓扑各Bolt具体实现,27个Class为编译产物,另有XML配置文件、MD说明文档和39张PNG流程截图,包体仅1.17MB,轻量完整,可快速对照源码和图示梳理运行过程。已有159人学习浏览,适合希望通过完整样例掌握Storm消费Kafka、规则匹配与多渠道告警的开发者,也可作为扩展规则引擎和告警模块的参考基础。

1. 日志监控告警系统.zip:一个能直接改的Storm+Kafka实时告警骨架

做运维监控的人常常会遇到一种尴尬:日志量不大时,用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 主题读日志 KafkaSpoutConfig<String, 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或单个.class,Ctrl+A导出全部源码。反编译出来的代码可能变量名变成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,再创建topic:

kafka-topics.sh --create \ --topic app_log \ --partitions 3 \ --replication-factor 1 \ --bootstrap-server localhost:9092

partitions设置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和timeout,ProcessDataBolt匹配后,会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混在一起。因为正则编译非常耗CPU,Bolt里每条消息都动态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。更简单的做法是把它写成一个静态Map:StormTickBolt只负责put,ProcessDataBolt只负责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返回fail,KafkaSpout重新发送同一条日志。解决方法是两个:第一条是发送成功后立即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持有自己的连接引用,如果整个拓扑并行度是4+2+2,20个连接足够了。connectionTimeout设3000毫秒,超过3秒直接报错,而不是无限等。我见过有人把连接超时设为60秒,结果数据库一慢,大量线程卡在getConnection上,拓扑假死。

5.3 用TopkeyCountBolt做告警频率聚合:从原始日志到TopN

TopkeyCountBolt这个类,从名字看是在做TopKey统计。它可以统计某一类异常在时间窗口内出现的次数,超过阈值就升级告警级别。这个Bolt的实现逻辑很简单,维护一个HashMap计数,然后在tick tuple触发时把TopN记录往下游发:

// TopkeyCountBolt 计数逻辑示意 private Map<String, Integer> countMap = new HashMap<>(); public void execute(Tuple tuple) { if (isTickTuple(tuple)) { // 输出当前窗口内次数最多的Top10异常key for (Map.Entry<String, 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日志→验证告警落地”这四步,希望帮到你。

本文还有配套的精品资源,点击获取

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

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

立即咨询