[TOC] #### 1. 前言 --- 你的系统上线了,有 10 个微服务在运行,某天线上出了问题,用户反馈「页面打不开」 你开始排查:查哪个服务的日志? + 订单服务 ?没问题 + 用户服务 ?没问题 + 支付服务 ?也没问题 查了半天,发现是网关服务的日志里有一行报错,但之前根本没看网关的日志 10 个服务,每个服务一台机器,日志分散在 10 台机器上,排查问题的时候,你得一台一台登上去看,效率极低 #### 2. 业务为什么难做 --- 日志分散在各个服务上,带来的问题不只是排查麻烦: 问题一:日志格式不统一 订单服务用 JSON 格式,用户服务用纯文本,支付服务用自定义格式,想做统一分析 ?先把格式统一了再说 问题二:日志量大,存储压力大 一个服务一天几个 G 的日志,10 个服务一天几十 G,直接存在本地磁盘,磁盘很快满了 问题三:实时分析困难 你想统计「最近 5 分钟有多少 error」,得去 10 台机器上分别 grep,然后汇总,做不到实时 #### 3. 不用 RabbitMQ 会怎样 --- 最简单的方式:每个服务把日志写到本地文件,然后用 Filebeat 之类的工具采集到 Elasticsearch ```plaintext 订单服务 → 本地日志 → Filebeat → Elasticsearch 用户服务 → 本地日志 → Filebeat → Elasticsearch 支付服务 → 本地日志 → Filebeat → Elasticsearch ``` 能用,但有几个问题: + 每个服务都要装 Filebeat,运维成本高 + Filebeat 和 Elasticsearch 直连,如果 ES 挂了,日志就丢了 + 不同服务的日志要发到不同的 ES 索引,配置复杂 #### 4. 用 RabbitMQ 怎么解决 --- 在中间加一层 RabbitMQ: ```plaintext 订单服务 → RabbitMQ → Elasticsearch 用户服务 → RabbitMQ → Elasticsearch 支付服务 → RabbitMQ → Elasticsearch ``` 所有服务把日志发到 RabbitMQ,消费者从 RabbitMQ 取日志,写入 Elasticsearch 服务和 ES 之间多了一层缓冲 #### 5. 代码实现 --- 服务端发送日志,封装一个简单的日志工具类: ```php class LogProducer { private $channel; public function __construct($channel) { $this->channel = $channel; } public function log($level, $module, $message, $context = []) { $data = [ 'level' => $level, 'module' => $module, 'message' => $message, 'context' => $context, 'timestamp' => date('Y-m-d H:i:s'), 'server' => gethostname(), ]; $amqpMessage = new AMQPMessage(json_encode($data), [ 'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT, ]); // 用 Topic Exchange,路由标识为 模块.级别 $this->channel->basic_publish($amqpMessage, 'log_exchange', "{$module}.{$level}"); } } // 使用 $logger = new LogProducer($channel); $logger->log('error', 'order', '数据库连接失败', ['host' => '10.0.0.1']); $logger->log('info', 'user', '用户登录成功', ['user_id' => 123]); ``` 声明 Exchange 和队列: ```php $channel->exchange_declare('log_exchange', 'topic', false, true, false); // error 日志 → 告警 $channel->queue_declare('alert_queue', false, true, false, false); $channel->queue_bind('alert_queue', 'log_exchange', '*.error'); // 全量日志 → ES $channel->queue_declare('es_queue', false, true, false, false); $channel->queue_bind('es_queue', 'log_exchange', '#'); ``` 消费者:写入 Elasticsearch ```php $callback = function ($msg) { $data = json_decode($msg->body, true); // 写入 ES Elasticsearch::index(['index' => 'logs-' . date('Y.m.d'), 'body' => $data]); $msg->ack(); }; $channel->basic_consume('es_queue', '', false, false, false, false, $callback); ``` 告警消费者 ```php $callback = function ($msg) { $data = json_decode($msg->body, true); // 发送告警 sendAlert("[{$data['module']}] {$data['message']}"); $msg->ack(); }; $channel->basic_consume('alert_queue', '', false, false, false, false, $callback); ``` 整个架构 ```plaintext 各个服务 ↓ 发日志 RabbitMQ (Topic Exchange) ├── *.error → alert_queue → 告警系统 ├── *.warning → monitor_queue → 监控系统 └── # → es_queue → Elasticsearch → Kibana ``` 不同级别的日志走不同的处理流程 加一个新服务 ?只要它往 RabbitMQ 发日志就行,不需要改消费者的代码 #### 6. 本文小结 --- RabbitMQ 在日志系统中的角色:日志中转站 | 能力 | 体现 | | ------------ | ------------ | | 解耦 | 服务和 ES 不直接连接 | | 削峰 | 日志量大的时候先存队列 | | 分流 | Topic Exchange 按级别分发 | | 可靠性 | 消息持久化,日志不丢 | 配合 ELK(Elasticsearch + Logstash + Kibana),就是一个完整的日志收集方案,下一讲,我们来做分布式事务