跳转到正文

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: correlatedpublisher-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。

队列与消费配置

配置默认值说明
两端 exchangeown.system-log持久化直连交换机
两端 routing-keysystem-log两端须一致
Producer persistence.enabledtrue是否注册发送保存器
Producer confirm-timeout5s等待发布确认的时限
Consumer enabledtrue消费与拓扑装配开关
Consumer auto-startuptrue是否自动启动监听容器
Consumer queueown.system-log接收队列
Consumer concurrency / prefetch1 / 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 的转发器,避免循环发送。采集端本地直写后又消费自身消息,仍可能重复保存到同一后端。