本文目录导读:

我来给你一个完整的Java RocketMQ案例,包括生产者和消费者的实现。
环境准备
Maven依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>4.9.4</version>
</dependency>
生产者示例
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class SimpleProducer {
public static void main(String[] args) throws Exception {
// 1. 创建生产者实例,指定生产者组名
DefaultMQProducer producer = new DefaultMQProducer("producer_group");
// 2. 指定NameServer地址
producer.setNamesrvAddr("localhost:9876");
// 3. 启动生产者
producer.start();
try {
// 4. 发送消息
for (int i = 0; i < 10; i++) {
// 创建消息,指定Topic、Tag和消息体
Message msg = new Message(
"TopicTest", // Topic
"TagA", // Tag
("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET) // 消息内容
);
// 5. 发送消息并获取结果
SendResult sendResult = producer.send(msg);
System.out.printf("发送结果: %s%n", sendResult);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
// 6. 关闭生产者
producer.shutdown();
}
}
}
消费者示例
1 简单消费者(Push模式)
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class SimpleConsumer {
public static void main(String[] args) throws Exception {
// 1. 创建消费者实例,指定消费者组名
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group");
// 2. 指定NameServer地址
consumer.setNamesrvAddr("localhost:9876");
// 3. 设置消费位置(从最新消费还是从头消费)
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// 4. 订阅Topic和Tag
consumer.subscribe("TopicTest", "*"); // "*"表示订阅所有Tag
// 5. 注册消息监听器
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
try {
for (MessageExt msg : msgs) {
String body = new String(msg.getBody(), "UTF-8");
System.out.printf("接收到消息: Topic=%s, Tag=%s, 消息内容=%s%n",
msg.getTopic(), msg.getTags(), body);
}
// 消费成功
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
// 消费失败,稍后重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
// 6. 启动消费者
consumer.start();
System.out.println("消费者启动成功");
}
}
2 有序消息消费者
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class OrderConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
// 订阅订单Topic
consumer.subscribe("OrderTopic", "*");
// 注册有序消息监听器
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeOrderlyContext context) {
context.setAutoCommit(true);
for (MessageExt msg : msgs) {
String orderId = msg.getKeys(); // 获取消息Key(订单ID)
String body = new String(msg.getBody());
System.out.printf("消费订单消息: 订单ID=%s, 内容=%s%n", orderId, body);
}
return ConsumeOrderlyStatus.SUCCESS;
}
});
consumer.start();
System.out.println("有序消费者启动成功");
}
}
有序消息生产者
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import java.util.List;
public class OrderProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("order_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
try {
// 模拟发送同一个订单的消息到同一个队列
for (int i = 0; i < 10; i++) {
// 假设订单ID为 order_001
String orderId = "order_001";
Message msg = new Message(
"OrderTopic",
"TagA",
orderId, // 设置消息Key为订单ID
("订单操作步骤 " + i).getBytes()
);
// 使用MessageQueueSelector确保同一个订单的消息发送到同一个队列
SendResult result = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
// 根据订单ID选择队列
String orderId = (String) arg;
long index = Math.abs(orderId.hashCode()) % mqs.size();
return mqs.get((int) index);
}
}, orderId);
System.out.printf("发送有序消息: %s%n", result);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.shutdown();
}
}
}
事务消息示例
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
public class TransactionProducer {
public static void main(String[] args) throws Exception {
TransactionMQProducer producer = new TransactionMQProducer("transaction_producer_group");
producer.setNamesrvAddr("localhost:9876");
// 设置事务监听器
producer.setTransactionListener(new TransactionListener() {
private ConcurrentHashMap<String, Boolean> localTrans = new ConcurrentHashMap<>();
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
try {
System.out.println("执行本地事务");
// 模拟本地事务操作(如数据库操作)
Thread.sleep(1000);
// 本地事务成功
localTrans.put(msg.getTransactionId(), true);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
// 本地事务失败,回滚
localTrans.put(msg.getTransactionId(), false);
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 检查本地事务状态
Boolean success = localTrans.get(msg.getTransactionId());
if (success != null && success) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.UNKNOW;
}
});
producer.start();
try {
Message msg = new Message("TransactionTopic", "TagA",
"事务消息测试".getBytes());
// 发送事务消息
SendResult result = producer.sendMessageInTransaction(msg, null);
System.out.printf("发送结果: %s%n", result);
} catch (Exception e) {
e.printStackTrace();
}
producer.shutdown();
}
}
批量消息示例
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.ArrayList;
import java.util.List;
public class BatchProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("batch_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
// 创建批量消息
List<Message> messages = new ArrayList<>();
for (int i = 0; i < 10; i++) {
Message msg = new Message("BatchTopic", "TagA",
("批量消息 " + i).getBytes());
messages.add(msg);
}
try {
// 发送批量消息
SendResult result = producer.send(messages);
System.out.printf("批量发送结果: %s%n", result);
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.shutdown();
}
}
}
运行说明
启动RocketMQ(需要先安装Docker)
# 使用Docker启动RocketMQ docker run -d -p 9876:9876 --name rocketmq-namesrv rocketmq-namesrv docker run -d -p 10911:10911 --name rocketmq-broker rocketmq-broker
运行步骤
- 先启动消费者程序
- 再启动生产者程序
- 观察消费者控制台输出的消息
注意事项
- 消息顺序性:需要有序消息时,必须使用
MessageQueueSelector确保消息路由到同一队列 - 消费重试:消费失败时返回
RECONSUME_LATER,消息会重试 - 消息去重:Producer发送消息时建议设置唯一Key,便于Consumer去重
- 异常处理:生产者和消费者都要妥善处理异常,避免资源泄漏
- 配置优化:生产环境需要根据实际需求调整生产者和消费者的配置参数
这个案例涵盖了RocketMQ的主要使用场景,你可以根据实际需求进行相应的调整和扩展。