StreamMQ 是一款基于 Redis Stream 与 Redisson 构建的开源消息中间件 SDK,以 MIT 协议发布。它将 Redis Stream 的原生能力封装为一套类 RocketMQ 的、面向业务开发者友好的消息 API,让你在无需引入重量级 MQ 集群的前提下,获得注解驱动消费、事务消息、延时消息、顺序消息等企业级特性。
已有 Redis?你已经拥有了消息中间件。StreamMQ 复用现有 Redis 基础设施,无需引入 NameServer、Broker、Zookeeper 等额外组件,一个 Redis 即是一个 MQ 集群。
对齐 RocketMQ RocketMQTemplate / @RocketMQMessageListener 的编程模型,迁移成本低,学习曲线平缓。如果你熟悉 RocketMQ,你已经会使用 StreamMQ。
事务消息、18 级延时消息、顺序消息、批量发送、死信队列、消息压缩、消息过滤——开箱即用的企业级能力,不输独立 MQ 集群。
自动装配、配置绑定、Actuator 端点、Micrometer 指标——与 Spring 生态无缝衔接,@EnableStreamMQ 一键开启。
序列化器、转换器、过滤器、拦截器、重试策略、重平衡策略、压缩编解码器、死信失败策略、管理鉴权器、链路追踪采集器——几乎一切可扩展。
651 个测试用例全部通过,覆盖核心消息能力、事务流程、延时投递、顺序消费、DLQ 处理等全场景。已在真实生产环境中验证。
┌─────────────────────────────────────────────────────────────────────────┐
│ StreamMQ Architecture │
├─────────────────────────────────────────────────────────────────────────┤
│ │
│ ┌───────────────────────────────────────────────────────────────────┐ │
│ │ Spring Boot Application │ │
│ │ ┌─────────────┐ ┌──────────────┐ ┌──────────────────────────┐ │ │
│ │ │@EnableStreamMQ│ │@StreamMQConsumer│ │ StreamMessageTemplate │ │ │
│ │ │ (自动装配) │ │ (声明式消费) │ │ (统一发送入口) │ │ │
│ │ └──────┬──────┘ └──────┬───────┘ └───────────┬──────────────┘ │ │
│ └─────────┼─────────────────┼─────────────────────┼────────────────┘ │
│ │ │ │ │
│ ┌─────────▼─────────────────▼─────────────────────▼────────────────┐ │
│ │ streammq-core │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │ │
│ │ │ Message │ │ Template │ │ Consumer │ │Producer │ │Transaction│ │
│ │ │ Builder │ │ Service │ │ Listener │ │ Factory │ │ Executor │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │ │
│ │ │ Filter │ │Interceptor│ │ Retry │ │ Rebalance│ │ DLQ │ │ │
│ │ │ Chain │ │ Chain │ │ Policy │ │ Strategy │ │Handler │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │ │
│ │ │Serializer│ │Converter │ │Codec │ │ Trace │ │Metrics │ │ │
│ │ │ (SPI) │ │ (SPI) │ │ (SPI) │ │Collector │ │ (Mic.)│ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │ │
│ └───────────────────────────────────────────────────────────────────┘ │
│ │ │ │ │
│ ┌─────────▼─────────────────▼─────────────────────▼────────────────┐ │
│ │ streammq-redisson │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │ │
│ │ │ Redisson │ │ Stream │ │ Delay │ │ PEL │ │ Tx │ │ │
│ │ │ Producer │ │ Listener │ │ Scheduler│ │ Claimer │ │Scanner │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │ │
│ └───────────────────────────────────────────────────────────────────┘ │
│ │ │
│ ┌─────────▼────────────────────────────────────────────────────────┐ │
│ │ Redis 7.2+ │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │ Stream │ │ ZSet │ │ Hash │ │ Sorted │ │ │
│ │ │ (消息存储)│ │(延时队列)│ │(事务状态)│ │ Set │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │ │
│ └──────────────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────────────┘
| 能力 | StreamMQ | Redisson RStream | Spring Data Redis Stream | RocketMQ | Kafka |
|---|---|---|---|---|---|
| 底层存储 | Redis Stream | Redis Stream | Redis Stream | NameServer+Broker | Broker+KRaft |
| 部署复杂度 | 低(仅 Redis) | 低(仅 Redis) | 低(仅 Redis) | 高(独立集群) | 高(独立集群) |
| 注解声明式消费 | 支持 | 不支持 | 部分支持 | 支持 | 不支持 |
| Template 编程模型 | 支持 | 不支持 | 不支持 | 支持 | 不支持 |
| 事务消息 | 支持 | 不支持 | 不支持 | 支持 | 不支持 |
| 延时消息 | 支持(18 级+任意) | 不支持 | 不支持 | 支持(18 级) | 不支持 |
| 顺序消息 | 支持 | 不支持 | 不支持 | 支持 | 支持(分区内) |
| 死信队列 | 支持(含二级 DLQ) | 不支持 | 不支持 | 支持 | 不支持 |
| 消息过滤 | Tag + SQL92 | 不支持 | 不支持 | 支持 | 不支持 |
| 消息压缩 | 支持(GZIP SPI) | 不支持 | 不支持 | 支持 | 支持 |
| 背压控制 | 支持(InflightQueue) | 不支持 | 不支持 | 支持 | 支持 |
| Spring Boot 3 集成 | 深度集成 | 一般 | 一般 | 一般(第三方) | 一般(第三方) |
| SPI 扩展点数量 | 12 个 | 0 | 0 | 少量 | 少量 |
| 管理接口 | REST API + Actuator | 无 | 无 | Dashboard | 无 |
| 链路追踪 | 支持(TraceCollector SPI) | 不支持 | 不支持 | 支持 | 不支持 |
| 学习成本 | 低 | 中 | 中 | 中 | 中 |
| 适用规模 | 中小规模(< 1 亿/天) | 中小规模 | 中小规模 | 大规模 | 超大规模 |
| 组件 | 最低版本 | 推荐版本 |
|---|---|---|
| JDK | 21 | 21+ |
| Maven | 3.9 | 3.9+ |
| Redis | 7.2 | 7.2+ |
| Spring Boot | 3.3 | 3.3.5 |
<dependencyManagement>
<dependencies>
<dependency>
<groupId>io.github.streammq</groupId>
<artifactId>streammq-bom</artifactId>
<version>0.1.0-SNAPSHOT</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>io.github.streammq</groupId>
<artifactId>streammq-spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson-spring-boot-starter</artifactId>
</dependency>
</dependencies>spring:
application:
name: streammq-demo
streammq:
enabled: true
namespace: streammq
redisson:
singleServerConfig:
address: "redis://127.0.0.1:6379"
database: 0@SpringBootApplication
@EnableStreamMQ
public class DemoApplication {
public static void main(String[] args) {
SpringApplication.run(DemoApplication.class, args);
}
}@Component
public class OrderService {
private final StreamMessageTemplate template;
public OrderService(StreamMessageTemplate template) {
this.template = template;
}
public SendResult sendOrder(String orderId, String content) {
Message<String> message = MessageBuilder.<String>withTopic("order-topic")
.tag("created")
.keys(orderId)
.body(content)
.userProperty("traceId", "t-001")
.build();
return template.syncSend(message);
}
}@Component
@StreamMQConsumer(topic = "order-topic", consumerGroup = "order-consumer-group")
public class OrderConsumer implements StreamMessageConcurrentlyConsumer<String> {
@Override
public ConsumeAction onMessage(Message<String> message, ConsumeContext context) {
System.out.println("收到订单:" + message.getKeys() + ", 内容:" + message.getBody());
return ConsumeAction.SUCCESS;
}
}就这样!启动应用,发送一条消息,消费者会自动接收并处理。
一行注解,声明式定义消费者,支持并发消费、顺序消费、广播消费、DLQ 消费四种模型。
// 并发消费(默认)
@StreamMQConsumer(topic = "order-topic", consumerGroup = "order-group")
// 顺序消费
@StreamMQConsumer(topic = "order-topic", consumerGroup = "order-group",
messageModel = MessageModel.ORDERLY, shardCount = 8)
// 广播消费
@StreamMQConsumer(topic = "order-topic", consumerGroup = "order-group",
consumeMode = ConsumeMode.BROADCASTING)
// DLQ 消费
@StreamMQConsumer(topic = "order-topic", consumerGroup = "order-group", dlqMode = true)统一的发送入口,支持同步、异步、单向、批量、事务五种发送方式。
public interface StreamMessageTemplate {
<T> SendResult syncSend(Message<T> message);
<T> SendResult syncSend(Message<T> message, long timeoutMillis);
<T> SendResult syncSend(Message<T> message, long timeoutMillis, int retryTimes);
<T> CompletableFuture<SendResult> asyncSend(Message<T> message);
<T> void asyncSend(Message<T> message, SendCallback callback);
<T> void sendOneway(Message<T> message);
<T> List<SendResult> syncSendBatch(BatchMessage<T> batch);
<T> SendResult executeInTransaction(Message<T> message, TransactionCallback<T> callback);
}半消息 + 本地事务 + 回查机制,保证最终一致性。
// 发送事务消息
TransactionCallback<String> callback = (message, ctx) -> {
try {
executeLocalTransaction(message.getBody());
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
};
SendResult result = template.executeInTransaction(message, callback);
// 事务回查器
@Component
@StreamMQTransactionConsumer(transactionGroup = "default-tx-group")
public class TransactionCheckerImpl implements TransactionChecker<String> {
@Override
public LocalTransactionState check(Message<String> message, TransactionContext context) {
return checkLocalTransactionStatus(context.getTransactionId());
}
}内置 18 级固定延时,亦可自定义任意毫秒延时。
// 固定延时(18 级)
Message<String> msg1 = MessageBuilder.<String>withTopic("delay-topic")
.body("content")
.delayLevel(DelayLevel.MINUTE_5) // 延时 5 分钟
.build();
// 任意延时毫秒
Message<String> msg2 = MessageBuilder.<String>withTopic("delay-topic")
.body("content")
.delayTimeMillis(15 * 60 * 1000L) // 延时 15 分钟
.build();延时级别对照表:
| 级别 | 延时 | 级别 | 延时 | 级别 | 延时 |
|---|---|---|---|---|---|
SECOND_1 |
1s | MINUTE_3 |
3m | MINUTE_20 |
20m |
SECOND_5 |
5s | MINUTE_4 |
4m | MINUTE_30 |
30m |
SECOND_10 |
10s | MINUTE_5 |
5m | HOUR_1 |
1h |
SECOND_30 |
30s | MINUTE_6 |
6m | HOUR_2 |
2h |
MINUTE_1 |
1m | MINUTE_7 |
7m | ||
MINUTE_2 |
2m | MINUTE_8 |
8m | ||
MINUTE_9 |
9m | MINUTE_10 |
10m |
基于 ShardingKey 的分片顺序消费,保证同一分区内严格有序。
@Component
@StreamMQConsumer(
topic = "order-topic",
consumerGroup = "order-orderly-group",
messageModel = MessageModel.ORDERLY,
shardCount = 8
)
public class OrderOrderlyConsumer implements StreamMessageOrderlyConsumer<String> {
@Override
public OrderlyAction onMessage(Message<String> message, ConsumeOrderlyContext context) {
processOrder(message.getBody());
return OrderlyAction.SUCCESS;
}
}
// 发送时指定 shardingKey
Message<String> message = MessageBuilder.<String>withTopic("order-topic")
.shardingKey("user-123")
.body("content")
.build();BatchMessage 批量投递,充分利用 Redis Pipeline 提升吞吐。
BatchMessage<String> batch = BatchMessage.<String>builder()
.topic("order-topic")
.addMessage(msg1)
.addMessage(msg2)
.addMessage(msg3)
.build();
List<SendResult> results = template.syncSendBatch(batch);消费重试耗尽后的消息自动进入 DLQ,支持二级 DLQ 与自定义失败策略。
@Component
@StreamMQConsumer(
topic = "order-topic",
consumerGroup = "order-group",
dlqMode = true
)
public class OrderDlqConsumer implements StreamMessageConcurrentlyConsumer<String> {
@Override
public ConsumeAction onMessage(Message<String> message, ConsumeContext context) {
handleDeadLetter(message);
return ConsumeAction.SUCCESS;
}
}支持 Tag 表达式与 SQL92 表达式两种过滤模式。
// Tag 过滤
@StreamMQConsumer(
topic = "order-topic",
consumerGroup = "order-group",
selectorExpression = "tag1 || tag2"
)
// SQL92 过滤
@StreamMQConsumer(
topic = "order-topic",
consumerGroup = "order-group",
selectorType = SelectorType.SQL92,
selectorExpression = "a = 1 AND b > 2"
)通过 CompressionCodec SPI 支持 GZIP 压缩,可配置压缩阈值。
streammq:
producer:
compress-threshold: 1024 # 消息体超过 1KB 时自动压缩
compression-codec: gzip # 压缩编解码器| 模块 | 说明 |
|---|---|
| streammq-bom | BOM(Bill of Materials),统一版本管理 |
| streammq-core | 核心抽象层,定义消息模型、API、SPI 接口(无 Spring 依赖) |
| streammq-redisson | Redisson 适配层,基于 Redis Stream 实现核心能力 |
| streammq-spring-boot-starter | Spring Boot 3 自动装配、配置绑定、Actuator 集成 |
| streammq-test | 测试工具包,提供嵌入式 Redis、断言工具、Mock 工具 |
| streammq-samples | 示例工程集合,覆盖快速开始、事务、延时、顺序、DLQ、拦截器 |
streammq:
enabled: true
namespace: streammq
# 生产者配置
producer:
group: default-producer-group
send-timeout: 3000
retry-times: 2
compress-threshold: 0 # 0=不压缩,>0 时超过阈值自动压缩
compression-codec: gzip
# 消费者配置
consumer:
consume-thread-min: 1
consume-thread-max: 64
pull-batch-size: 32
max-reconsume-times: 16
consume-timeout: 30000
inflight-capacity: 1000 # 背压队列容量
pull-interval: 0 # 拉取间隔(毫秒)
# 事务配置
transaction:
check-interval-ms: 60000 # 回查间隔
max-check-times: 15 # 最大回查次数
batch-size: 32 # 单次扫描批量
# 死信队列配置
dlq:
enabled: true
max-retry-times: 3
# 可观测性配置
metrics:
enabled: true
tracing:
enabled: false
# 管理接口配置
management:
enabled: true| 属性 | 类型 | 默认值 | 说明 |
|---|---|---|---|
topic |
String | - | 主题(必填) |
consumerGroup |
String | - | 消费组(必填) |
messageModel |
MessageModel | CONCURRENT | 消费模型:CONCURRENT / ORDERLY |
consumeMode |
ConsumeMode | CLUSTERING | 消费模式:CLUSTERING / BROADCASTING |
consumeThreadMin |
int | 1 | 最小消费线程数 |
consumeThreadMax |
int | 64 | 最大消费线程数 |
maxReconsumeTimes |
int | 16 | 最大重试次数 |
consumeTimeout |
long | 30000 | 消费超时(毫秒) |
pullBatchSize |
int | 32 | 单次拉取批量 |
selectorExpression |
String | "*" | Tag/SQL92 过滤表达式 |
selectorType |
SelectorType | TAG | 过滤类型:TAG / SQL92 |
shardCount |
int | 4 | 顺序消费分片数 |
dlqMode |
boolean | false | 是否 DLQ 消费者 |
pullInterval |
long | 0 | 拉取间隔(毫秒) |
streamMaxLen |
int | 0 | Stream 最大长度(0=不限制) |
retryStreamMaxLen |
int | 0 | 重试 Stream 最大长度 |
enableMsgTrace |
boolean | false | 是否启用消息追踪 |
serializer |
Class | MessageSerializer.class | 序列化器(默认全局) |
messageConverter |
Class | MessageConverter.class | 消息转换器(默认全局) |
retryPolicy |
Class | RetryPolicy.class | 重试策略(默认全局) |
rebalanceStrategy |
Class | RebalanceStrategy.class | 重平衡策略(默认全局) |
dlqFailureStrategy |
Class | DlqFailureStrategy.class | 死信失败策略(默认全局) |
consumerFilter |
Class[] | {} | 消费者专属过滤器 |
完整配置参考请查看 配置文档。
StreamMQ 通过 SPI 提供丰富的扩展点,几乎一切可替换:
| SPI 接口 | 作用 | 默认实现 |
|---|---|---|
MessageSerializer |
消息序列化/反序列化 | JacksonJsonSerializer / JdkSerializer |
MessageConverter |
消息体与业务对象转换 | DefaultMessageConverter / CompactMessageConverter / PassThroughMessageConverter |
ProducerFilter |
生产者过滤器(过滤链) | 无默认 |
ConsumerFilter |
消费者过滤器(全局+per-consumer) | TagSelectorFilter / SqlSelectorFilter |
ProducerInterceptor |
生产者拦截器(拦截链) | 无默认 |
ConsumerInterceptor |
消费者拦截器(拦截链) | 无默认 |
RetryPolicy |
重试策略 | FixedArrayRetryPolicy |
RebalanceStrategy |
消费者重平衡策略 | AverageRebalanceStrategy / ConsistentHashRebalanceStrategy |
CompressionCodec |
消息压缩编解码 | GzipCompressionCodec |
TraceCollector |
链路追踪上下文采集 | Slf4jTraceCollector |
ManagementAuthenticator |
管理接口鉴权 | AllowAllAuthenticator / BasicAuthAuthenticator / TokenAuthenticator / DenyAllAuthenticator |
DlqFailureStrategy |
死信消费失败策略 | AbstractDlqFailureStrategy |
@Component
public class CustomMessageSerializer implements MessageSerializer {
@Override
public byte[] serialize(Object obj) throws SerializationException {
// 自定义序列化逻辑
return customSerialize(obj);
}
@Override
public <T> T deserialize(byte[] bytes, Class<T> type) throws SerializationException {
// 自定义反序列化逻辑
return customDeserialize(bytes, type);
}
@Override
public String name() {
return "custom";
}
}// 在注解中指定使用自定义 SPI
@StreamMQConsumer(
topic = "order-topic",
consumerGroup = "order-group",
serializer = CustomMessageSerializer.class,
consumerFilter = { CustomFilter.class }
)| 指标名 | 类型 | 说明 |
|---|---|---|
streammq.producer.send.total |
Counter | 发送消息总数 |
streammq.producer.send.success |
Counter | 发送成功数 |
streammq.producer.send.failed |
Counter | 发送失败数 |
streammq.producer.send.duration |
Timer | 发送耗时分布 |
streammq.consumer.consume.total |
Counter | 消费消息总数 |
streammq.consumer.consume.duration |
Timer | 消费耗时分布 |
streammq.consumer.retry.total |
Counter | 重试消息数 |
streammq.consumer.dlq.total |
Counter | 进入死信队列数 |
| 端点 | 说明 |
|---|---|
/actuator/health |
健康检查(含 StreamMQ 组件状态) |
/actuator/metrics |
Micrometer 指标 |
/actuator/prometheus |
Prometheus 格式指标 |
| 端点 | 方法 | 说明 |
|---|---|---|
/streammq/admin/consumer-groups |
GET | 查询消费组列表 |
/streammq/admin/consumer-groups/{group} |
GET | 查询消费组详情 |
/streammq/admin/topics |
GET | 查询 Topic 列表 |
/streammq/admin/topics/{topic} |
GET | 查询 Topic 详情 |
/streammq/admin/dlq/{group} |
GET | 查询 DLQ 消息 |
/streammq/admin/dlq/{group} |
DELETE | 清空 DLQ |
/streammq/admin/ack |
POST | 手动 ACK |
/streammq/admin/rebalance |
POST | 触发重平衡 |
通过 TraceCollector SPI 采集链路上下文,支持 traceId 透传。
MDC.put("traceId", "t-001");
template.syncSend(message); // traceId 自动透传到消费者| 示例 | 说明 |
|---|---|
| streammq-sample-quickstart | 快速开始示例 |
| streammq-sample-transaction | 事务消息示例 |
| streammq-sample-delay | 延时消息示例 |
| streammq-sample-orderly | 顺序消息示例 |
| streammq-sample-dlq | 死信队列示例 |
| streammq-sample-interceptor | 拦截器示例 |
| 文档 | 说明 |
|---|---|
| 首页 | 项目首页 |
| 快速开始 | 5 分钟上手指南 |
| 核心特性 | 完整特性文档 |
| 核心概念 | 关键术语解释 |
| API 文档 | 完整 API 参考 |
| 配置参考 | 全部配置项 |
| 部署指南 | 生产环境部署 |
| FAQ | 常见问题解答 |
| 贡献指南 | 参与贡献 |
| 文档 | 说明 |
|---|---|
| 产品需求文档 | V1.0 PRD |
| 架构设计 | V1.0 架构设计 |
| 功能设计 | V1.0 功能设计 |
| 详细设计 | V1.0 详细设计 |
| 文档 | 说明 |
|---|---|
| V2.0 PRD | V2.0 产品需求文档 |
| V2.0 架构设计 | V2.0 架构设计 |
| V2.0 功能设计 | V2.0 功能设计 |
- 注解驱动消费(
@StreamMQConsumer) -
StreamMessageTemplate编程模型(同步/异步/单向/批量/事务) - 集群消费 + 广播消费
- 顺序消费(ShardingKey 分片)
- 事务消息(半消息 + 回查)
- 延时消息(18 级 + 任意毫秒)
- 死信队列(含二级 DLQ)
- 消息过滤(Tag + SQL92)
- 消息压缩(GZIP)
- 背压控制(InflightQueue)
- 消费超时自动取消
- Micrometer 指标 + MDC 日志
- 链路追踪(TraceCollector SPI)
- 管理 REST API
- 12 个 SPI 扩展点
- Spring Boot 3 自动装配 + Actuator 集成
- 多后端抽象层(BackendProvider SPI,支持 Redis / Kafka / RabbitMQ / Pulsar)
- Kafka 后端实现(基于 Kafka Client 的 BackendProvider)
- 跨机房复制(异步复制,RPO ≤ 1s)
- Kafka 线网协议兼容(原生 Kafka Client 零代码接入)
- Spring Cloud Stream Binder(实现 Spring Cloud Stream Binder SPI)
- Kubernetes Operator(CRD + Operator,弹性伸缩、配置热更新)
- 消息画像与拓扑图(可视化消息流转拓扑)
- 分布式追踪增强(OpenTelemetry 集成)
欢迎参与 StreamMQ 开源建设!请阅读 贡献指南 了解详细信息。
# 1. Fork & Clone
git clone https://github.com/<your-username>/streammq.git
cd streammq
# 2. 创建分支
git checkout -b feature/your-feature
# 3. 编写代码 & 测试
mvn clean test
# 4. 提交(遵循 Conventional Commits)
git commit -m "feat: add your feature"
# 5. 发起 PR- Bug 报告:提交 Issue,描述问题与复现步骤
- 功能请求:提交 Issue,描述期望功能与使用场景
- 代码贡献:提交 Pull Request,关联相关 Issue
- 文档改进:完善文档、修正错误、补充示例
- 问题解答:在 Discussions 中帮助其他用户
- GitHub Issues:https://github.com/streammq/streammq/issues
- GitHub Discussions:https://github.com/streammq/streammq/discussions
- Pull Requests:https://github.com/streammq/streammq/pulls
- 已有 Redis 基础设施,希望复用为消息总线
- 中小规模业务(单集群日消息量 < 1 亿)
- 需要事务消息 / 延时消息 / 顺序消息能力但不想引入独立 MQ 集群
- 微服务架构下基于 Spring Boot 3 的轻量级异步通信
- 电商订单状态流转、支付回调、库存扣减、通知推送
- 超大规模流式数据处理(单集群日消息量 > 1 亿)—— 建议使用 Kafka
- 对消息吞吐要求极高且可容忍少量丢失 —— 建议使用 Kafka
- 需要复杂路由规则(topic 通配符、多级路由)—— 建议使用 RabbitMQ
- 已有成熟 MQ 集群且无 Redis 资源 —— 直接复用现有 MQ
| 技术 | 版本 | 用途 |
|---|---|---|
| Java | 21+ | 运行时 |
| Spring Boot | 3.3.5 | 框架基础 |
| Redisson | 3.34.1 | Redis 客户端 |
| Jackson | 2.18.1 | JSON 序列化 |
| Fury | 0.9.0 | 高性能序列化(可选) |
| Protostuff | 1.8.0 | Protobuf 序列化(可选) |
| Lombok | - | 代码简化 |
| Micrometer | - | 指标收集 |
| SLF4J | - | 日志门面 |
本项目基于 MIT License 开源。