Kafka项目简介

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可视化平台 → 监控面板展示+指标趋势分析

分层说明

  1. 数据采集层:Filebeat采集日志、node-exporter采集系统指标
  2. 数据传输层:Kafka集群承载日志流传输
  3. 数据处理层:Python消费者解析日志、Celery调度告警任务
  4. 数据存储层:MySQL存储结构化日志、Prometheus存储时序指标
  5. 告警通知层:双渠道异步告警,保障异常及时触达
  6. 可视化层: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

快速部署流程

  1. 基础环境配置
    关闭防火墙、SELinux,配置主机名与静态IP,安装基础依赖工具
  2. Kafka三节点集群部署(KRaft模式)
    生成集群UUID→格式化存储目录→配置server.properties→启动集群→验证可用性
  3. 中间件部署
    部署MySQL(创建日志库表)、Redis(Celery代理)
  4. 日志采集服务部署
    三节点部署Filebeat,配置日志采集路径与Kafka输出
  5. 核心业务服务部署
    部署Python日志消费服务→部署Celery告警服务→部署Flask测试应用
  6. 监控可视化部署
    部署node-exporter→部署Prometheus→部署Grafana→配置数据源与监控面板
  7. 服务自启配置
    所有核心服务编写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. 可视化验证

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部署,适配企业级生产环境

项目亮点

  1. 无ZK依赖:KRaft模式简化架构,降低部署维护成本
  2. 全开源免费:无商业软件授权,降低企业投入
  3. 轻量化低侵入:资源占用低,不影响Kafka核心业务
  4. 全链路闭环:日志采集到告警可视化一站式解决
  5. 高可用稳定:服务自启、自动重启、异常重试机制完善
  6. 易扩展适配:支持新增节点、新增告警渠道、功能快速扩展