外观
RabbitMQ 集中保存
采集应用将日志发送到 RabbitMQ,存储应用接收后写入数据库或文件。两端可以独立部署。
text
采集应用 → RabbitMQ → 存储应用 → 数据库 / 文件本例使用 CSV 作为消费端存储。开始前,两端应用应已通过 Business BOM 管理版本,并准备好可连接的 RabbitMQ。
1. 配置采集应用
添加生产模块:
xml
<dependency>
<groupId>com.own.business</groupId>
<artifactId>springboot-business-system-log-persistence-mq-rabbit-producer</artifactId>
</dependency>开启采集和发布确认:
yaml
spring:
application:
name: account-service
rabbitmq:
host: localhost
port: 5672
username: ${RABBITMQ_USERNAME}
password: ${RABBITMQ_PASSWORD}
connection-timeout: 5s
publisher-confirm-type: correlated
publisher-returns: true
own:
system-log:
enabled: true
mq-rabbit-producer:
exchange: own.system-log
routing-key: system-log
confirm-timeout: 5s按快速开始标注需要记录的接口。Producer 自动注册为保存器,接收采集日志并发布消息。publisher-confirm-type: correlated 和 publisher-returns: true 都必须配置,否则保存器装配失败。
2. 配置存储应用
添加消费模块和实际存储模块:
xml
<dependencies>
<dependency>
<groupId>com.own.business</groupId>
<artifactId>springboot-business-system-log-persistence-mq-rabbit-consumer</artifactId>
</dependency>
<dependency>
<groupId>com.own.business</groupId>
<artifactId>springboot-business-system-log-persistence-file-csv</artifactId>
</dependency>
</dependencies>配置队列和 CSV 目录:
yaml
spring:
rabbitmq:
host: localhost
port: 5672
username: ${RABBITMQ_USERNAME}
password: ${RABBITMQ_PASSWORD}
connection-timeout: 5s
own:
system-log:
enabled: false
mq-rabbit-consumer:
enabled: true
exchange: own.system-log
routing-key: system-log
queue: own.system-log
concurrency: 1
prefetch: 32
file-csv:
save-path: ./logs/system-log/csv
persistence: true这里关闭的是存储应用自身的接口采集。MQ 消费独立运行,直接调用保存器,因此本例 CSV 的 own.system-log.file-csv.persistence 必须保持 true;数据库目标使用各自的 persistence.enabled。没有实际保存器时消费应用启动失败。
3. 启动并验证
先启动存储应用,再启动采集应用。消费者会声明交换机、队列和绑定;生产者只声明交换机,没有接收队列时会产生不可路由失败。
调用采集应用中标注的接口,等待消费完成后查看存储应用的 CSV。记录应保留采集端的事件 ID、服务名、发起人、起止时间和快照状态。
要增加集中保存目标,只需在存储应用引入其他存储模块、配置连接或文件目录并启用对应保存器。例如增加 JSON 文件存储 后,同一条消息可同时写入 CSV 和完整日志 JSON Lines。
队列与消费配置
| 配置 | 默认值 | 说明 |
|---|---|---|
两端 exchange | own.system-log | 持久化直连交换机 |
两端 routing-key | system-log | 两端须一致 |
Producer persistence.enabled | true | 是否注册发送保存器 |
Producer confirm-timeout | 5s | 等待发布确认的时限 |
Consumer enabled | true | 消费与拓扑装配开关 |
Consumer auto-startup | true | 是否自动启动监听容器 |
Consumer queue | own.system-log | 接收队列 |
Consumer concurrency / prefetch | 1 / 32 | 并发数 / 每消费者预取数 |
Producer 前缀为 own.system-log.mq-rabbit-producer,Consumer 前缀为 own.system-log.mq-rabbit-consumer。
同一 queue 的多个实例竞争消费。需要每个应用收到完整副本时,使用不同 queue,绑定同一 exchange/routing-key。
保存失败如何处理
消费者在监听线程中依次调用所有实际保存器,全部成功后才确认消息。单个后端失败仍尝试其他后端,最终只要存在失败,就拒绝该消息并进入死信队列。
死信交换机为 <exchange>.dead,死信队列为 <queue>.dead。失败不自动重入队。修复故障后再安排死信重放,并考虑已成功后端的重复写入;JSON、CSV 和 SQL 文件不会自动去重。
生产端入队成功不等于已发布,Broker 确认也不等于下游已保存。生产端不提供自动重试或持久化 outbox,确认超时可能意味着投递结果不确定。
消息约定
每条消息为 UTF-8 JSON,正文是完整 SystemLog,AMQP type 为 own.system-log.v1,messageId 为原始事件 ID,最大 64 KiB。非法协议、JSON、缺失必要字段或超限消息会被拒绝。
消费者不再次调用采集记录器,也会排除实现 SystemLogForwardingPersistence 的转发器,避免循环发送。采集端本地直写后又消费自身消息,仍可能重复保存到同一后端。