Kafka集群分布式日志收集与实时监控告警平台
基于 Kafka(KRaft无ZK模式) + Filebeat + Python + Celery + MySQL + Prometheus + Grafana 构建的分布式日志采集、结构化存储、实时异常告警、系统指标可视化一体化运维监控平台,专为中小型企业Kafka集群运维场景设计
项目简介
随着分布式系统广泛应用,Kafka作为高吞吐消息队列已成为企业核心组件,但其集群运维面临监控盲区、告警滞后、人工巡检低效、日志难以追溯等问题。
本项目基于全开源技术栈,构建日志采集→传输→解析→存储→告警→可视化全链路闭环方案,实现Kafka三节点集群的全方位可观测性保障。
平台具备轻量化、低侵入、高可用、易扩展特性,不影响Kafka核心业务运行,可快速部署落地,有效提升集群运维效率与故障响应速度。
核心功能
1. Kafka集群高可用部署
- 采用KRaft模式部署3节点Kafka集群,彻底移除ZooKeeper依赖
- 支持分区副本机制,保障消息高可用传输
- 完成集群可用性、消息收发、副本同步全功能验证
2. 多节点日志实时采集
- 三节点统一部署Filebeat轻量级采集器
- 实时采集Kafka服务日志、Flask应用访问/错误日志
- 支持断点续传、日志文件轮转、批量压缩发送
- 采集延迟≤10秒,无日志丢失与重复采集
3. 日志可靠传输
- 以Kafka
nginx-log主题作为日志缓冲通道 - 3分区对应3节点日志源,提升并行传输效率
- 依托Kafka高可用特性,单节点故障不中断日志传输
4. 日志结构化解析与存储
- Python多进程消费者程序,正则表达式解析非结构化日志
- 全量访问日志存入
nginx_access_logs表 - 5xx服务端错误日志单独存入
error_logs表 - 支持按时间、状态码、IP、URL多维度快速查询
- 数据库连接异常自动重连,保障数据写入稳定
5. 双渠道实时异步告警
- 基于Celery+Redis实现异步告警任务调度
- 4xx客户端错误:触发QQ邮箱单渠道告警
- 5xx服务端错误:触发QQ邮箱+钉钉双渠道告警
- 支持告警失败自动重试(最多3次)、防告警风暴
- 告警内容包含节点、时间、IP、URL、错误详情等关键信息
6. Flask测试验证载体
- 开发轻量级Web测试服务,提供正常访问接口
- 提供
/trigger-500-error接口主动触发500错误 - 访问不存在路径自动生成404错误日志
- 用于全链路日志采集、存储、告警有效性验证
7. 系统指标可视化监控
- Prometheus+node-exporter采集三节点CPU、内存、磁盘、网络指标
- Grafana搭建专业监控面板,支持秒级数据刷新
- 单节点指标展示、多节点对比、趋势分析
- 指标阈值可视化提醒,直观掌握集群运行状态
8. 服务高可用保障
- 所有核心服务基于systemd托管
- 支持开机自启、进程异常自动重启
- 服务依赖顺序配置,避免启动失败
- 运行日志持久化,便于故障排查
系统架构
本平台采用分层架构设计,各模块解耦、协同工作,整体架构如下:
用户访问请求
│
▼
Flask测试应用(kafka1/kafka2/kafka3)→ 生成标准化访问/错误日志
│
▼
Filebeat采集器(各节点)→ 实时捕获日志数据
│
▼
Kafka 3节点集群(KRaft模式)→ 日志主题:nginx-log
│
▼
Python多进程消费者 → 正则解析日志→结构化处理
│
├─→ MySQL数据库 → 全量访问日志+5xx错误日志分类存储
│
└─→ Celery异步任务队列(Redis代理)→ 告警规则判断
│
├─→ QQ邮箱告警(4xx/5xx错误)
└─→ 钉钉群机器人告警(5xx严重错误)
│
▼
Prometheus指标采集 → 节点系统硬件指标收集
│
▼
Grafana可视化平台 → 监控面板展示+指标趋势分析
分层说明
- 数据采集层:Filebeat采集日志、node-exporter采集系统指标
- 数据传输层:Kafka集群承载日志流传输
- 数据处理层:Python消费者解析日志、Celery调度告警任务
- 数据存储层:MySQL存储结构化日志、Prometheus存储时序指标
- 告警通知层:双渠道异步告警,保障异常及时触达
- 可视化层:Grafana面板呈现监控数据,支撑运维决策
技术栈明细
| 技术组件 | 版本 | 用途 |
|---|---|---|
| 操作系统 | CentOS 7.9 | 集群运行环境 |
| 消息队列 | Kafka 3.3.1 | 日志缓冲传输,KRaft无ZK模式 |
| 日志采集 | Filebeat 7.17.0 | 轻量级日志采集、推送至Kafka |
| 开发语言 | Python 3.9.16 | 日志消费、告警、Flask应用开发 |
| 异步任务 | Celery 5.2.7 + Redis 6.2.7 | 异步告警任务调度、队列代理 |
| 关系数据库 | MySQL 8.0.33 | 结构化日志存储、查询 |
| 监控采集 | Prometheus 2.45.0 | 系统指标时序存储、查询 |
| 指标采集 | node-exporter 1.6.1 | 节点CPU/内存/磁盘/网络指标采集 |
| 可视化 | Grafana 9.5.2 | 监控面板搭建、数据展示 |
| Web框架 | Flask 2.3.3 | 测试载体Web应用开发 |
| 服务管理 | systemd | 服务自启、故障自动恢复 |
| 告警通道 | QQ邮箱(SMTP)、钉钉机器人(Webhook) | 双渠道异常通知 |
部署环境要求
硬件配置(3节点统一)
- CPU:4核
- 内存:8GB
- 磁盘:100GB
- 网络:三节点内网互通,关闭防火墙与SELinux
节点分工
| 节点IP | 节点名称 | 部署组件 |
|---|---|---|
| 192.168.126.130 | kafka1 | Kafka、Filebeat、MySQL、Redis、Prometheus、Grafana、消费者、告警服务、Flask |
| 192.168.126.131 | kafka2 | Kafka、Filebeat、node-exporter、Flask |
| 192.168.126.132 | kafka3 | Kafka、Filebeat、node-exporter、Flask |
快速部署流程
- 基础环境配置
关闭防火墙、SELinux,配置主机名与静态IP,安装基础依赖工具 - Kafka三节点集群部署(KRaft模式)
生成集群UUID→格式化存储目录→配置server.properties→启动集群→验证可用性 - 中间件部署
部署MySQL(创建日志库表)、Redis(Celery代理) - 日志采集服务部署
三节点部署Filebeat,配置日志采集路径与Kafka输出 - 核心业务服务部署
部署Python日志消费服务→部署Celery告警服务→部署Flask测试应用 - 监控可视化部署
部署node-exporter→部署Prometheus→部署Grafana→配置数据源与监控面板 - 服务自启配置
所有核心服务编写systemd配置文件,开启开机自启
功能验证方法
1. 日志采集验证
# 生成测试日志
curl http://节点IP:5000
curl http://节点IP:5000/test-404
curl http://节点IP:5000/trigger-500-error
# 查看Kafka日志接收
/opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server 192.168.126.130:9092 --topic nginx-log --from-beginning
2. 日志存储验证
mysql -u root -p密码 -D log_analysis
SELECT server_name,client_ip,request_url,status_code FROM nginx_access_logs ORDER BY id DESC LIMIT 10;
SELECT * FROM error_logs ORDER BY id DESC LIMIT 10;
3. 告警功能验证
- 触发404错误→等待1分钟→查看QQ邮箱告警
- 触发500错误→等待1分钟→查看QQ+钉钉双渠道告警
4. 可视化验证
- Prometheus:http://192.168.126.130:9090
- Grafana:http://192.168.126.130:3000(默认admin/admin)
5. 服务自启验证
# 重启服务器
reboot
# 查看所有核心服务状态
systemctl status kafka filebeat mysqld redis log-consumer celery-worker celery-beat flask-web prometheus grafana-server
核心性能指标
经系统测试,平台核心性能完全满足设计要求:
- 日志采集延迟:≤10秒
- 告警响应时间:≤30秒
- 可视化面板刷新延迟:≤5秒
- 单节点CPU使用率:≤10%
- 单节点内存占用:≤2GB
- 对Kafka集群业务性能影响:可忽略
服务管理命令
# 启动服务
systemctl start 服务名
# 停止服务
systemctl stop 服务名
# 重启服务
systemctl restart 服务名
# 查看服务状态
systemctl status 服务名
# 开启开机自启
systemctl enable 服务名
# 关闭开机自启
systemctl disable 服务名
未来演进规划
v2.0 功能扩展
- 接入Kafka Exporter,实现Kafka业务指标深度监控
- 新增告警分级、告警处理状态、告警静默功能
- 支持自定义告警阈值与高流量访问告警
v3.0 技术优化
- 集成ELK Stack,实现双链路日志检索与分析
- 深化Grafana可视化,增加日志趋势、告警统计面板
- 优化日志解析逻辑,支持更多日志格式
v4.0 容器化改造
- 全组件Docker容器化封装
- 提供Docker Compose一键编排部署
- 支持快速迁移、弹性扩缩容
v5.0 智能化升级
- 集成大模型API,实现自然语言日志查询
- 智能故障诊断与自动化处理建议
- 支持K8s部署,适配企业级生产环境
项目亮点
- 无ZK依赖:KRaft模式简化架构,降低部署维护成本
- 全开源免费:无商业软件授权,降低企业投入
- 轻量化低侵入:资源占用低,不影响Kafka核心业务
- 全链路闭环:日志采集到告警可视化一站式解决
- 高可用稳定:服务自启、自动重启、异常重试机制完善
- 易扩展适配:支持新增节点、新增告警渠道、功能快速扩展