scaffold-message 异步消息模块
模块概述
scaffold-message 是 Scaffold v2 平台的异步消息通知模块,采用 SPI 双层架构设计:
- Transport 层(传输层):对接 MQ 中间件,负责消息的可靠传输。支持 Kafka、RabbitMQ、InMemory 三种实现,通过 classpath 自动检测激活。
- Channel 层(渠道层):定义消息的投递方式。内置 HttpCallback、InAppNotification、Log 三种渠道,支持自定义扩展。
业务模块注入 MessageBus 即可发送消息,无需关心底层 MQ 中间件和渠道实现。
功能列表
- 消息发送:通过 MessageBus 统一 API 发送异步消息
- 消息持久化:每条消息发送前落库(t_message_record),确保可追溯
- 消息重试:失败消息通过指数退避策略自动重试
- 消息重发:支持手动重发失败消息
- 通知中心:应用内通知的发送、查询、已读标记
- 消息记录管理:分页查询、状态跟踪
核心组件
核心类
| 类 |
说明 |
| MessageBus |
消息总线,统一入口 API。提供 send() 异步发送和 sendAndDeliver() 同步投递 |
| Message |
消息信封,包含 channel、topic、payload、headers 等 |
| MessageResult |
发送/投递结果 |
| MessageStatus |
消息状态枚举 |
实体
| 实体 |
表名 |
说明 |
| MessageRecord |
t_message_record |
消息记录(持久化) |
| MessageNotification |
t_message_notification |
应用内通知 |
服务层
| 服务 |
说明 |
| MessageRecordService |
消息记录 CRUD、状态更新、重发 |
| NotificationService |
通知发送、查询、已读标记 |
重试机制
| 类 |
说明 |
| RetryPolicy |
重试策略接口 |
| ExponentialBackoffRetry |
指数退避重试实现 |
| MessageRecordRetryScheduler |
定时扫描失败消息并重试 |
配置参数
配置前缀:scaffold.message
| 参数 |
类型 |
默认值 |
说明 |
| enabled |
boolean |
true |
模块开关 |
| transport |
String |
auto |
传输层选择:auto / kafka / rabbitmq / memory |
| retry.maxAttempts |
int |
3 |
最大重试次数 |
| retry.backoffMs |
String |
"1000,2000,4000" |
退避间隔(毫秒),逗号分隔,如 1s -> 2s -> 4s |
| threadPool.coreSize |
int |
4 |
线程池核心线程数 |
| threadPool.maxSize |
int |
8 |
线程池最大线程数 |
| threadPool.queueCapacity |
int |
1000 |
线程池队列容量 |
| threadPool.threadNamePrefix |
String |
"message-" |
线程名前缀 |
| kafka.topic |
String |
scaffold-message |
Kafka topic 名称 |
| kafka.groupId |
String |
scaffold-message-group |
Kafka 消费者 group-id |
| rabbit.exchange |
String |
scaffold.message.exchange |
RabbitMQ TopicExchange 名称 |
| rabbit.queue |
String |
scaffold.message.queue |
RabbitMQ 队列名称 |
| rabbit.routingPattern |
String |
# |
RabbitMQ 路由通配符 |
API 接口列表
消息记录管理(/v1/message/record)
| 方法 |
路径 |
说明 |
| POST |
/page |
分页查询消息记录 |
| POST |
/resend/{id} |
重发消息 |
通知中心(/v1/message/notification)
| 方法 |
路径 |
说明 |
| POST |
/page |
分页查询通知列表 |
| POST |
/send |
发送通知 |
| POST |
/read/{id} |
标记单条已读 |
| POST |
/read-all |
标记全部已读 |
| GET |
/unread-count |
获取未读数 |
SPI 扩展点
Transport 层(传输层)
传输层通过 MessageTransport 接口扩展。激活策略:
| 实现类 |
激活条件 |
配置类 |
| KafkaTransport |
classpath 存在 spring-kafka |
KafkaTransportAutoConfiguration |
| RabbitTransport |
classpath 存在 spring-amqp |
RabbitTransportAutoConfiguration |
| InMemoryTransport |
上述都不满足时降级 |
InMemoryTransportAutoConfiguration |
自定义传输层:实现 MessageTransport 接口并注册为 Spring Bean。
Channel 层(渠道层)
渠道层通过 MessageChannel 接口扩展。内置实现:
| 实现类 |
渠道名称 |
说明 |
| HttpCallbackChannel |
callback |
HTTP 回调通知 |
| InAppNotificationChannel |
notification |
应用内通知 |
| LogChannel |
log |
日志输出(开发调试用) |
自定义渠道:实现 MessageChannel 接口(getChannelName()、supports()、deliver()),注册为 Spring Bean 即可被 ChannelRouter 自动发现。
使用示例
1. 注入 MessageBus 发送消息
@RequiredArgsConstructor
@Service
public class OrderService {
private final MessageBus messageBus;
public void onOrderPaid(Order order) {
// 异步发送:经 MQ 传输 -> 渠道投递
messageBus.send("callback", "pay.notify",
"{\"orderId\":\"" + order.getId() + "\",\"status\":\"PAID\"}");
// 同步投递:不走 MQ,直接路由到渠道
Message msg = Message.builder()
.channel("notification")
.topic("order.paid")
.payload("{\"orderId\":\"" + order.getId() + "\"}")
.build();
messageBus.sendAndDeliver(msg);
}
}
2. 发送应用内通知
// 通过 NotificationSendDTO 发送给指定用户
NotificationSendDTO dto = new NotificationSendDTO();
dto.setUserId(userId);
dto.setTitle("订单支付成功");
dto.setContent("订单 " + orderId + " 已支付");
notificationService.send(dto);
3. 添加自定义渠道
@Component
public class EmailChannel implements MessageChannel {
@Override
public String getChannelName() { return "email"; }
@Override
public boolean supports(String channelName) { return "email".equals(channelName); }
@Override
public MessageResult deliver(Message message) {
// 实现邮件发送逻辑
return MessageResult.ok(message.getId());
}
}
消息状态流转
消息记录(MessageRecord)的状态值:
| 状态值 |
含义 |
| 0 |
PENDING - 待发送(已落库) |
| 1 |
DELIVERING - 传输中(已发送到 MQ) |
| 2 |
SUCCESS - 投递成功 |
| 3 |
FAILED - 投递失败(等待重试) |
注意:MessageStatus 为多态状态机,不遵循布尔语义规范。
注意事项
transport=auto 时按 classpath 自动选择 MQ 中间件,优先级:Kafka > RabbitMQ > InMemory
- InMemoryTransport 仅适用于开发环境,不保证消息持久化和可靠性
- 消息重试通过
MessageRecordRetryScheduler 定时扫描失败记录并重新投递
MessageBus.send() 为异步模式:消息先落库 PENDING -> 传输层发送 -> 更新 DELIVERING -> 渠道消费后回写 SUCCESS/FAILED
MessageBus.sendAndDeliver() 为同步模式:不走 MQ,直接路由到渠道并同步返回结果