若依微服务消息中心架构设计与优化实践 1. 项目背景与需求分析在若依(RuoYi)微服务架构中消息通知模块最初采用SysNotice表结构进行设计随着业务复杂度提升这种单一表结构的设计逐渐暴露出以下问题扩展性不足SysNotice表结构固定难以支持多种消息类型系统通知、待办提醒、预警消息等渠道耦合严重消息发送逻辑与业务代码深度耦合新增短信、邮件等渠道需要修改多处代码状态管理混乱消息的已读/未读状态、处理状态等缺乏统一管理机制性能瓶颈高并发场景下频繁操作数据库表容易成为系统瓶颈2. 架构设计与技术选型2.1 整体架构设计统一消息中心采用分层架构设计[接入层] ├── REST API ├── RPC接口 └── 消息队列监听 [核心服务层] ├── 消息路由引擎 ├── 模板管理 ├── 渠道适配器 └── 状态机管理 [存储层] ├── MySQL结构化数据 ├── MongoDB非结构化内容 └── Redis缓存与队列2.2 关键技术选型Spring Cloud Stream用于解耦消息生产与消费支持RabbitMQ/Kafka等消息中间件状态模式实现消息生命周期管理创建→发送→接收→处理→归档模板引擎采用Freemarker实现消息内容动态渲染分布式事务使用Seata保证消息创建与业务操作的一致性3. 核心实现细节3.1 数据模型重构// 消息主体 Entity Table(name sys_message) public class SysMessage { Id GeneratedValue(strategy IDENTITY) private Long id; Enumerated(STRING) private MessageType type; // 消息类型 private String templateCode; // 模板编码 private String title; private String content; Enumerated(STRING) private MessageStatus status; Column(updatable false) private LocalDateTime createTime; } // 消息接收者 Entity Table(name sys_message_receiver) public class MessageReceiver { Id GeneratedValue(strategy IDENTITY) private Long id; private Long messageId; private Long userId; Enumerated(STRING) private ReadStatus readStatus; private LocalDateTime readTime; }3.2 消息发送流程优化public interface MessageSender { SendResult send(MessageDTO message); } Service public class MessageSendService { Autowired private ListMessageSender senders; Transactional public void sendMessage(MessageDTO message) { // 1. 持久化消息 SysMessage entity convertToEntity(message); messageRepository.save(entity); // 2. 选择发送渠道 MessageSender sender senders.stream() .filter(s - s.support(message.getChannel())) .findFirst() .orElseThrow(); // 3. 异步发送 messageQueue.push(() - { SendResult result sender.send(message); updateMessageStatus(entity.getId(), result); }); } }3.3 消息模板管理采用模板与内容分离的设计# 消息模板示例 (yaml配置) templates: - code: TASK_ASSIGN title: 任务分配提醒 content: | 您有新的任务待处理 任务名称${taskName} 优先级${priority} 截止时间${deadline?string(yyyy-MM-dd)} channels: [WEB,EMAIL]4. 与原有系统集成方案4.1 兼容性处理数据迁移编写Flyway迁移脚本将SysNotice数据转换到新结构INSERT INTO sys_message (id, type, title, content, status, create_time) SELECT id, SYSTEM, title, content, CASE status WHEN 1 THEN UNREAD ELSE READ END, create_time FROM sys_notice;API适配层提供兼容原有接口的FacadeDeprecated RestController RequestMapping(/notice) public class NoticeCompatController { Autowired private MessageService messageService; GetMapping(/list) public R list(SysNotice notice) { // 将旧参数转换为新参数 MessageQuery query convertQuery(notice); return R.ok(messageService.queryList(query)); } }4.2 灰度发布策略通过Nacos配置中心控制新旧系统切换RestController RequestMapping(/api/notice) public class NoticeRouterController { Value(${notice.new.enabled:false}) private boolean useNewSystem; Autowired private OldNoticeService oldService; Autowired private MessageService newService; GetMapping(/{id}) public R detail(PathVariable Long id) { return useNewSystem ? newService.getDetail(id) : oldService.getNotice(id); } }5. 性能优化实践5.1 读写分离设计Repository public class MessageRepositoryImpl implements MessageRepository { Autowired Qualifier(messageWriteMapper) private MessageMapper writeMapper; Autowired Qualifier(messageReadMapper) private MessageMapper readMapper; Override Transactional(readOnly true) public PageMessageVO queryPage(MessageQuery query) { return readMapper.selectPage(query); } Override Transactional public void save(SysMessage message) { if (message.getId() null) { writeMapper.insert(message); } else { writeMapper.update(message); } } }5.2 缓存策略采用多级缓存架构热点数据Caffeine本地缓存有效期5分钟普通数据Redis集群过期时间30分钟持久层MySQL分库分表Service CacheConfig(cacheNames message) public class MessageServiceImpl implements MessageService { Cacheable(key detail: #id, unless #result null) Override public MessageDetail getDetail(Long id) { return mapper.selectDetailById(id); } Caching(evict { CacheEvict(key detail: #message.id), CacheEvict(key list: #message.userId) }) Override public void updateStatus(MessageStatusUpdateDTO dto) { // 更新逻辑 } }6. 监控与运维方案6.1 监控指标埋点Aspect Component RequiredArgsConstructor public class MessageMonitorAspect { private final MeterRegistry meterRegistry; Around(execution(* com.ruoyi.message..*.*(..))) public Object monitor(ProceedingJoinPoint pjp) throws Throwable { String method pjp.getSignature().getName(); Timer.Sample sample Timer.start(meterRegistry); try { Object result pjp.proceed(); sample.stop(meterRegistry.timer(message.process, method, method, status, success)); return result; } catch (Exception e) { sample.stop(meterRegistry.timer(message.process, method, method, status, fail)); throw e; } } }6.2 告警规则配置在Prometheus中配置关键指标告警groups: - name: message-center rules: - alert: HighFailureRate expr: rate(message_process_fail_total[1m]) / rate(message_process_total[1m]) 0.05 for: 5m labels: severity: warning annotations: summary: 消息处理失败率过高 description: 当前失败率 {{ $value }}7. 迁移实施经验在实际迁移过程中我们总结了以下关键经验双写过渡期新旧系统并行运行2周通过对比日志确保数据一致性批量处理优化使用MyBatis的批量插入代替单条插入性能提升8倍动态开关在Nacos中配置功能开关出现问题时快速回滚补偿机制编写定时任务修复状态不一致的消息记录Scheduled(cron 0 0/5 * * * ?) public void fixInconsistentStatus() { // 查询处理中超时的消息 ListLong timeoutMessages mapper.selectTimeoutMessages( LocalDateTime.now().minusMinutes(30)); timeoutMessages.forEach(id - { try { MessageDetail detail getDetail(id); if (detail.getStatus() SENDING) { log.warn(Fix timeout message: {}, id); updateStatus(id, FAILED); } } catch (Exception e) { log.error(Fix message failed: id, e); } }); }

本月热点