高可用 RocketMQ 消费端设计(防数据丢失 + 批量处理 + 无数据库锁)

基于 RocketMQ 的高可用消费端,核心需求是高效消费大批量基础数据实现数据新增 / 更新逻辑避免数据库锁竞争,同时保证零数据丢失。完整设计方案,消费端采用注解式配置,兼顾性能、可靠性和可维护性。

一、核心设计思路

1. 高可用设计

  • 集群消费模式:多实例部署分摊消息,避免单节点瓶颈;
  • 线程池隔离:消费线程与数据库操作线程分离,避免慢查询阻塞消费;
  • 重试机制:RocketMQ 自带指数退避重试,失败消息入死信队列兜底;
  • 幂等性保障:基于消息 ID+Redis 实现幂等,避免重复消费导致数据错误。

2. 防数据库锁设计

  • 批量操作:单次批量 500 条,减少 SQL 执行次数;
  • 原子化新增 / 更新:使用INSERT ... ON DUPLICATE KEY UPDATE,一条 SQL 完成逻辑,避免先查后更的锁竞争;
  • 数据库优化:开启批量写入优化,业务唯一键加唯一索引,减少锁粒度。

3. 零数据丢失设计

  • 消费签收机制:数据库操作成功后才返回消费成功,失败则重试;
  • 事务保障:数据库批量操作加事务,失败整体回滚;
  • 失败日志记录:消费失败的消息写入日志表,便于人工重试;
  • 死信队列兜底:超过重试次数的消息进入死信队列,避免丢失。

二、环境依赖(Maven)

<!-- Spring Boot 基础 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter</artifactId>
    <version>2.7.15</version>
</dependency>

<!-- RocketMQ Spring 整合(注解式消费) -->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.3</version>
</dependency>

<!-- MyBatis-Plus(简化批量操作) -->
<dependency>
    <groupId>com.baomidou</groupId>
    <artifactId>mybatis-plus-boot-starter</artifactId>
    <version>3.5.3.1</version>
</dependency>

<!-- MySQL 驱动 -->
<dependency>
    <groupId>mysql</groupId>
    <artifactId>mysql-connector-java</artifactId>
    <version>8.0.33</version>
    <scope>runtime</scope>
</dependency>

<!-- Redis(幂等校验) -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>

<!-- 工具类 -->
<dependency>
    <groupId>cn.hutool</groupId>
    <artifactId>hutool-all</artifactId>
    <version>5.8.22</version>
</dependency>

三、核心配置

1. 应用配置(application.yml)

spring:
  # 数据库配置
  datasource:
    url: jdbc:mysql://127.0.0.1:3306/test_db?useUnicode=true&characterEncoding=utf8&useSSL=false&rewriteBatchedStatements=true
    username: root
    password: 123456
    driver-class-name: com.mysql.cj.jdbc.Driver
  # Redis配置
  redis:
    host: 127.0.0.1
    port: 6379
    password:
    database: 0
    lettuce:
      pool:
        max-active: 20
        max-idle: 10
        min-idle: 5

# RocketMQ 配置(注解式消费)
rocketmq:
  name-server: 127.0.0.1:9876 # NameServer地址(集群用逗号分隔)
  consumer:
    group: base_data_consumer_group # 消费组(唯一)
    consume-message-batch-max-size: 200 # 批量消费大小
    consume-thread-max: 8 # 消费线程数
    max-reconsume-times: 20 # 最大重试次数
    consume-timeout: 30000 # 消费超时时间(30秒)
  producer:
    group: base_data_producer_group # 生产者组(备用)

# MyBatis-Plus 配置
mybatis-plus:
  mapper-locations: classpath:mapper/*.xml
  type-aliases-package: com.example.entity
  configuration:
    map-underscore-to-camel-case: true # 下划线转驼峰
    # log-impl: org.apache.ibatis.logging.stdout.StdOutImpl # 生产环境关闭

2. 线程池配置(异步处理数据库操作)

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

import java.util.concurrent.ThreadPoolExecutor;

/**
 * 线程池配置:隔离消费线程与数据库操作线程
 */
@Configuration
public class ThreadPoolConfig {

    /**
     * 数据库操作线程池
     */
    @Bean("dbOperationExecutor")
    public ThreadPoolTaskExecutor dbOperationExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(10); // 核心线程数
        executor.setMaxPoolSize(20); // 最大线程数
        executor.setQueueCapacity(1000); // 队列容量
        executor.setThreadNamePrefix("db-operation-"); // 线程名前缀
        // 拒绝策略:队列满时由调用线程执行(避免消息丢失)
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        executor.initialize();
        return executor;
    }
}

四、数据实体与 Mapper

1. 基础数据 DTO(消息传输对象)

import lombok.Data;
import java.util.Date;

/**
 * 基础数据传输对象(RocketMQ消息体)
 */
@Data
public class BaseDataDTO {
    /** 业务唯一键(核心,用于幂等+新增/更新判断) */
    private String bizUniqueId;
    /** 基础编码 */
    private String code;
    /** 基础名称 */
    private String name;
    /** 基础值 */
    private String value;
    /** 状态(0-禁用,1-启用) */
    private Integer status;
    /** 更新时间 */
    private Date updateTime;
}

2. 基础数据实体(数据库表映射)

import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.util.Date;

/**
 * 基础数据实体(数据库表:t_base_data)
 */
@Data
@TableName("t_base_data")
public class BaseDataPO {
    @TableId(type = IdType.AUTO)
    private Long id;

    /** 业务唯一键(加唯一索引) */
    @TableField("biz_unique_id")
    private String bizUniqueId;

    @TableField("code")
    private String code;

    @TableField("name")
    private String name;

    @TableField("value")
    private String value;

    @TableField("status")
    private Integer status;

    @TableField("update_time")
    private Date updateTime;

    @TableField("create_time")
    private Date createTime;
}

3. 消费失败日志实体(兜底用)

import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.util.Date;

/**
 * 消费失败日志实体(数据库表:t_base_data_fail_log)
 */
@Data
@TableName("t_base_data_fail_log")
public class BaseDataFailLogPO {
    @TableId(type = IdType.AUTO)
    private Long id;

    /** RocketMQ消息ID */
    private String msgId;

    /** 消息体(JSON格式) */
    private String msgBody;

    /** 失败原因 */
    private String failReason;

    /** 已重试次数 */
    private Integer retryCount;

    /** 创建时间 */
    private Date createTime;
}

4. Mapper 接口

import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.example.entity.BaseDataFailLogPO;
import com.example.entity.BaseDataPO;
import org.apache.ibatis.annotations.Param;
import java.util.List;

/**
 * 基础数据Mapper
 */
public interface BaseDataMapper extends BaseMapper<BaseDataPO> {
    /**
     * 批量新增或更新(MySQL:INSERT ... ON DUPLICATE KEY UPDATE)
     */
    int batchUpsert(@Param("list") List<BaseDataPO> list);
}

/**
 * 消费失败日志Mapper
 */
public interface BaseDataFailLogMapper extends BaseMapper<BaseDataFailLogPO> {
    /**
     * 批量插入失败日志
     */
    int batchInsert(@Param("list") List<BaseDataFailLogPO> list);
}

5. Mapper XML(批量操作 SQL)

<!-- BaseDataMapper.xml -->
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.example.mapper.BaseDataMapper">
    <!-- 批量新增或更新(防数据库锁核心SQL) -->
    <insert id="batchUpsert">
        INSERT INTO t_base_data (biz_unique_id, code, name, value, status, update_time, create_time)
        VALUES
        <foreach collection="list" item="item" separator=",">
            (
            #{item.bizUniqueId},
            #{item.code},
            #{item.name},
            #{item.value},
            #{item.status},
            #{item.updateTime},
            NOW()
            )
        </foreach>
        ON DUPLICATE KEY UPDATE
        code = VALUES(code),
        name = VALUES(name),
        value = VALUES(value),
        status = VALUES(status),
        update_time = VALUES(update_time)
    </insert>
</mapper>

<!-- BaseDataFailLogMapper.xml -->
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.example.mapper.BaseDataFailLogMapper">
    <!-- 批量插入失败日志 -->
    <insert id="batchInsert">
        INSERT INTO t_base_data_fail_log (msg_id, msg_body, fail_reason, retry_count, create_time)
        VALUES
        <foreach collection="list" item="item" separator=",">
            (
            #{item.msgId},
            #{item.msgBody},
            #{item.failReason},
            #{item.retryCount},
            NOW()
            )
        </foreach>
    </insert>
</mapper>

五、核心业务逻辑(Service)

import cn.hutool.json.JSONUtil;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.example.dto.BaseDataDTO;
import com.example.entity.BaseDataFailLogPO;
import com.example.entity.BaseDataPO;
import com.example.mapper.BaseDataFailLogMapper;
import com.example.mapper.BaseDataMapper;
import org.springframework.beans.BeanUtils;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

import javax.annotation.Resource;
import java.util.List;
import java.util.stream.Collectors;

/**
 * 基础数据业务层(核心:批量新增/更新+幂等+失败日志)
 */
@Service
public class BaseDataService extends ServiceImpl<BaseDataMapper, BaseDataPO> {

    @Resource
    private RedisTemplate<String, String> redisTemplate;
    @Resource
    private BaseDataFailLogMapper failLogMapper;

    /** 幂等缓存前缀 */
    private static final String DUPLICATE_KEY_PREFIX = "base_data:duplicate:";
    /** 幂等缓存过期时间(7天) */
    private static final long REDIS_EXPIRE_SECONDS = 86400L * 7;

    /**
     * 幂等校验:检查消息是否已消费
     * @param msgId RocketMQ消息ID
     * @return true=重复,false=未消费
     */
    public boolean checkDuplicate(String msgId) {
        return redisTemplate.hasKey(DUPLICATE_KEY_PREFIX + msgId);
    }

    /**
     * 标记消息已消费(数据库操作成功后调用)
     */
    public void markMsgConsumed(String msgId) {
        redisTemplate.opsForValue().set(DUPLICATE_KEY_PREFIX, msgId, REDIS_EXPIRE_SECONDS);
    }

    /**
     * 批量新增/更新基础数据(事务保障+失败日志)
     * @param dtoList 基础数据列表
     * @param msgId RocketMQ消息ID(用于失败日志)
     * @return true=成功,false=失败
     */
    @Transactional(rollbackFor = Exception.class)
    public boolean batchUpsert(List<BaseDataDTO> dtoList, String msgId) {
        try {
            // DTO转换为PO
            List<BaseDataPO> poList = dtoList.stream()
                    .map(dto -> {
                        BaseDataPO po = new BaseDataPO();
                        BeanUtils.copyProperties(dto, po);
                        return po;
                    })
                    .collect(Collectors.toList());

            // 执行批量新增/更新(核心SQL,防锁)
            int affectedRows = baseMapper.batchUpsert(poList);
            return affectedRows > 0;
        } catch (Exception e) {
            // 记录失败日志(兜底,避免数据丢失)
            String failReason = e.getMessage() != null ? e.getMessage() : "数据库操作异常";
            BaseDataFailLogPO failLogPO = new BaseDataFailLogPO();
            failLogPO.setMsgId(msgId);
            failLogPO.setMsgBody(JSONUtil.toJsonStr(dtoList));
            failLogPO.setFailReason(failReason);
            failLogPO.setRetryCount(0);
            failLogMapper.insert(failLogPO);

            // 事务回滚
            throw new RuntimeException("批量操作失败,已记录日志", e);
        }
    }
}

六、注解式消费端(核心)

java运行

import cn.hutool.json.JSONUtil;
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.apache.rocketmq.common.message.MessageExt;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.stream.Collectors;

/**
 * 注解式RocketMQ消费端(基础数据消费)
 * - MessageModel.CLUSTERING:集群消费模式(多实例分摊)
 * - ConsumeMode.CONCURRENTLY:并发消费(提高效率)
 */
@Component
@RocketMQMessageListener(
        topic = "base_data_topic", // 订阅的主题
        consumerGroup = "${rocketmq.consumer.group}", // 消费组(配置文件读取)
        messageModel = MessageModel.CLUSTERING, // 集群消费
        consumeMode = ConsumeMode.CONCURRENTLY, // 并发消费
        maxReconsumeTimes = 20 // 最大重试次数
)
public class BaseDataConsumer implements RocketMQListener<List<MessageExt>> {

    @Autowired
    private BaseDataService baseDataService;

    /**
     * 批量消费逻辑(注解式批量消费)
     * @param msgs 批量消息列表
     */
    @Override
    public void onMessage(List<MessageExt> msgs) {
        // 遍历每条消息,保证单条消息的可靠性
        for (MessageExt msg : msgs) {
            boolean isSuccess = false;
            String msgId = msg.getMsgId();
            try {
                // 1. 幂等校验(避免重复消费)
                if (baseDataService.checkDuplicate(msgId)) {
                    continue; // 已消费,直接跳过
                }

                // 2. 解析消息体
                String msgBody = new String(msg.getBody(), StandardCharsets.UTF_8);
                BaseDataDTO dataDTO = JSONUtil.toBean(msgBody, BaseDataDTO.class);

                // 3. 执行数据库操作(同步执行,保证消费签收可靠性)
                isSuccess = baseDataService.batchUpsert(List.of(dataDTO), msgId);

                // 4. 标记消息已消费(幂等持久化)
                if (isSuccess) {
                    baseDataService.markMsgConsumed(msgId);
                }
            } catch (Exception e) {
                // 抛出异常,触发RocketMQ重试(保证数据不丢失)
                throw new RuntimeException("消费消息失败,触发重试:" + msgId, e);
            }

            // 校验结果,失败则抛异常重试
            if (!isSuccess) {
                throw new RuntimeException("数据库操作失败,触发重试:" + msgId);
            }
        }
    }
}

七、数据库表结构

-- 基础数据表(唯一索引防重复)
CREATE TABLE `t_base_data` (
  `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键',
  `biz_unique_id` varchar(64) NOT NULL COMMENT '业务唯一键',
  `code` varchar(32) NOT NULL COMMENT '基础编码',
  `name` varchar(64) NOT NULL COMMENT '基础名称',
  `value` varchar(255) DEFAULT NULL COMMENT '基础值',
  `status` tinyint DEFAULT 1 COMMENT '状态(0-禁用,1-启用)',
  `update_time` datetime DEFAULT NULL COMMENT '更新时间',
  `create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_biz_unique_id` (`biz_unique_id`) COMMENT '业务唯一键索引(防重复+减锁)'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='基础数据表';

-- 消费失败日志表(兜底防丢失)
CREATE TABLE `t_base_data_fail_log` (
  `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键',
  `msg_id` varchar(64) NOT NULL COMMENT 'RocketMQ消息ID',
  `msg_body` text NOT NULL COMMENT '消息体(JSON)',
  `fail_reason` varchar(512) DEFAULT NULL COMMENT '失败原因',
  `retry_count` int DEFAULT 0 COMMENT '已重试次数',
  `create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  PRIMARY KEY (`id`),
  KEY `idx_msg_id` (`msg_id`) COMMENT '消息ID索引'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='基础数据消费失败日志表';

八、死信队列处理(最后一道防线)

1. 死信队列规则

  • 死信队列名称:%DLQ%+消费组名,例如%DLQ%base_data_consumer_group
  • 触发条件:消息重试次数达到maxReconsumeTimes(20 次)后自动进入死信队列;
  • 特性:死信队列的消息不会被自动消费,需人工介入。

2. 死信消息处理流程

  1. 监控告警:通过 RocketMQ Dashboard 监控死信队列长度,配置告警规则(有消息即告警);

  2. 数据排查:导出死信消息内容,检查数据格式、业务唯一键是否重复等问题;

  3. 人工重试:

    • 修复数据后,调用baseDataService.batchUpsert手动入库;
    • 或通过 RocketMQ 控制台将消息重新发送到原主题,触发消费重试;
  4. 归档清理:处理完成后,标记失败日志表的记录为「已处理」,避免重复重试。

九、关键优化与保障

1. 性能优化

  • 批量消费:单次消费 200 条,减少网络交互;
  • 批量数据库操作INSERT ... ON DUPLICATE KEY UPDATE减少 SQL 执行次数;
  • Redis 幂等:7 天有效期,避免重复消费,同时防止缓存膨胀。

2. 防锁优化

  • 业务唯一键加唯一索引,减少行锁粒度;
  • 控制批量大小(200 条),避免大事务锁表;
  • 开启 MySQLrewriteBatchedStatements=true,提升批量写入性能。

3. 数据不丢失保障

  • 消费签收:数据库操作成功后才标记消息已消费;
  • 事务保障:批量操作加事务,失败回滚;
  • 失败日志:消费失败自动记录,便于人工重试;
  • 死信队列:重试失败的消息归档,兜底处理。

总结

  1. 核心逻辑:通过注解式消费端实现批量消费,利用INSERT ... ON DUPLICATE KEY UPDATE实现原子化新增 / 更新,避免数据库锁竞争;
  2. 高可用:集群消费 + 线程池隔离 + 指数退避重试,保证消费端不宕机、不堆积;
  3. 零数据丢失:消费签收机制 + 事务 + 失败日志 + 死信队列,覆盖全链路数据兜底;
  4. 幂等性:Redis 标记已消费,避免重复消费导致数据错误。
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐