package com.mq.rocket.consumer;
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.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
/**
* 消息接收者
*/
public class Consumer {
public static void main(String[] args) throws MQClientException {
// 1.创建消费者consumer,指定消费者组名
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group1");
// 2.指定Nameserver地址
consumer.setNamesrvAddr("127.0.0.1:9876");
// 设置超时
consumer.setConsumeTimeout(15000);
// 3.订阅主题topic和tag
consumer.subscribe("topic1", "tag1");
// 4.设置回调函数,处理消息
consumer.registerMessageListener(new MessageListenerConcurrently() {
// 接收消息内容
@Override
public ConsumeConcurrentlyStatus consumeMessage(List msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
for (MessageExt msg:msgList) {
System.out.println(new String(msg.getBody()));
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 5.启动消费者consumer
consumer.start();
}
}
二、消费模式
1.负载均衡模式(默认):消费者共同消费
MessageModel.CLUSTERING
发送10条消息
消费者1:消费了3条
消费者2:消费了5条
消费者3:消费了2条
2.广播模式:每个消费者都消费同样的消息
MessageModel.BROADCASTING
package com.mq.rocket.consumer;
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.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import java.util.List;
/**
* 消息接收者
*/
public class Consumer {
public static void main(String[] args) throws MQClientException {
// 1.创建消费者consumer,指定消费者组名
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group1");
// 2.指定Nameserver地址
consumer.setNamesrvAddr("127.0.0.1:9876");
// 设置超时
consumer.setConsumeTimeout(15000);
// 3.订阅主题topic和tag
consumer.subscribe("topic1", "tag1");
// 设置消费模式:负载均衡和广播模式,默认负载均衡模式-MessageModel.CLUSTERING ,广播模式-MessageModel.BROADCASTING
consumer.setMessageModel(MessageModel.BROADCASTING);
// 4.设置回调函数,处理消息
consumer.registerMessageListener(new MessageListenerConcurrently() {
// 接收消息内容
@Override
public ConsumeConcurrentlyStatus consumeMessage(List msgList, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
for (MessageExt msg:msgList) {
System.out.println(new String(msg.getBody()));
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 5.启动消费者consumer
consumer.start();
}
}
欢迎分享,转载请注明来源:内存溢出
评论列表(0条)