高可用 RocketMQ 消费端设计(防数据丢失 + 批量处理 + 无数据库锁)
·
高可用 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. 死信消息处理流程
-
监控告警:通过 RocketMQ Dashboard 监控死信队列长度,配置告警规则(有消息即告警);
-
数据排查:导出死信消息内容,检查数据格式、业务唯一键是否重复等问题;
-
人工重试:
- 修复数据后,调用
baseDataService.batchUpsert手动入库; - 或通过 RocketMQ 控制台将消息重新发送到原主题,触发消费重试;
- 修复数据后,调用
-
归档清理:处理完成后,标记失败日志表的记录为「已处理」,避免重复重试。
九、关键优化与保障
1. 性能优化
- 批量消费:单次消费 200 条,减少网络交互;
- 批量数据库操作:
INSERT ... ON DUPLICATE KEY UPDATE减少 SQL 执行次数; - Redis 幂等:7 天有效期,避免重复消费,同时防止缓存膨胀。
2. 防锁优化
- 业务唯一键加唯一索引,减少行锁粒度;
- 控制批量大小(200 条),避免大事务锁表;
- 开启 MySQL
rewriteBatchedStatements=true,提升批量写入性能。
3. 数据不丢失保障
- 消费签收:数据库操作成功后才标记消息已消费;
- 事务保障:批量操作加事务,失败回滚;
- 失败日志:消费失败自动记录,便于人工重试;
- 死信队列:重试失败的消息归档,兜底处理。
总结
- 核心逻辑:通过注解式消费端实现批量消费,利用
INSERT ... ON DUPLICATE KEY UPDATE实现原子化新增 / 更新,避免数据库锁竞争; - 高可用:集群消费 + 线程池隔离 + 指数退避重试,保证消费端不宕机、不堆积;
- 零数据丢失:消费签收机制 + 事务 + 失败日志 + 死信队列,覆盖全链路数据兜底;
- 幂等性:Redis 标记已消费,避免重复消费导致数据错误。
更多推荐



所有评论(0)