使用 DefaultMQPushConsumer 的 registerMessageListener 方法注册监听器接收rocketmq消息
·
在RocketMQ中,DefaultMQPushConsumer 是一个用于接收消息的消费者实现。你可以使用 registerMessageListener 方法来注册一个消息监听器,以便在接收到消息时执行某些操作。
下面是一个使用 DefaultMQPushConsumer 的 registerMessageListener 方法来注册消息监听器的简单示例:
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.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class RocketMQConsumer {
public static void main(String[] args) throws MQClientException {
// 1. 实例化消费者
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("your_consumer_group");
// 2. 设置NameServer的地址
consumer.setNamesrvAddr("127.0.0.1:9876");
// 3. 订阅一个或多个Topic,以及Tag来过滤需要消费的消息
consumer.subscribe("your_topic", "*");
// 4. 注册回调实现类来处理从broker拉取回来的消息
consumer.registerMessageListener((List<MessageExt> msgs, ConsumeConcurrentlyContext context) -> {
for (MessageExt msg : msgs) {
System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);
}
// 标记该消息已经被成功消费
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
// 5. 启动消费者实例
consumer.start();
System.out.println("Consumer Started.");
}
}
在上面的代码中,我们首先创建了一个 DefaultMQPushConsumer 的实例,并设置了NameServer的地址。然后,我们订阅了一个或多个Topic(以及Tag来过滤消息)。
接下来,我们使用 registerMessageListener 方法注册了一个消息监听器。这个监听器是一个lambda表达式,它接受一个消息列表和一个上下文对象作为参数。在这个lambda表达式中,我们简单地打印出接收到的消息,并返回 ConsumeConcurrentlyStatus.CONSUME_SUCCESS 来标记消息已经被成功消费。
最后,我们调用 consumer.start() 方法来启动消费者实例。
注意:
- 你需要根据你的实际情况替换
"your_consumer_group"、"your_topic"和NameServer的地址"127.0.0.1:9876"。 - 在实际的生产环境中,你可能需要处理更复杂的消息处理逻辑,例如错误处理、重试策略等。
- 你需要确保你的RocketMQ的NameServer和Broker已经在运行,并且你的消费者能够连接到它们。
更多推荐
所有评论(0)