RocketMQ
简介
RocketMQ 是阿里巴巴开源的分布式消息中间件,具有高可用、高性能、低延迟的特点。支持分布式事务消息、顺序消息、定时消息等高级特性。
核心概念
- Producer:消息生产者,发送消息
- Consumer:消息消费者,接收消息
- Topic:主题,消息的第一级分类
- Tag:标签,消息的第二级分类,用于精细化过滤
- Message Queue:消息队列,Topic 的最小存储单元
- Broker:消息存储和转发服务器
消息模型
Producer -> NameServer -> Broker -> Consumer- NameServer:无状态的路由发现中心,提供 Topic 路由信息
- Broker:消息存储节点,主从架构保证高可用
消息类型
java
// 普通消息
DefaultMQProducer producer = new DefaultMQProducer("producer_group");
producer.start();
Message msg = new Message("TopicTest", "TagA", "Hello RocketMQ".getBytes());
SendResult result = producer.send(msg);
// 顺序消息
producer.send(msg, (mqs, msgObj, arg) -> {
int index = (int) arg % mqs.size();
return mqs.get(index);
}, orderId);
// 事务消息
TransactionMQProducer transactionProducer = new TransactionMQProducer("group");
transactionProducer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
return LocalTransactionState.COMMIT_MESSAGE;
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 事务回查
return LocalTransactionState.COMMIT_MESSAGE;
}
});适用场景
- 异步解耦:订单系统与库存、物流系统解耦
- 流量削峰:秒杀活动中的高并发流量缓冲
- 日志处理:分布式系统的日志采集与分发
- 分布式事务:基于 RocketMQ 的事务消息实现最终一致性