DefaultMQProducer producer = new DefaultMQProducer("ProducerGroup");
producer.setNamesrvAddr("localhost:9876");
producer.start();
Message message = new Message("TopicTest",
"TagA",
"Hello RocketMQ".getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
}
@Override
public void onException(Throwable e) {
}
});
producer.shutdown();
DefaultMQProducer producer = new DefaultMQProducer("ProducerGroup");
producer.setNamesrvAddr("localhost:9876");
producer.start();
List<Message> messages = new ArrayList<>();
messages.add(new Message("TopicTest", "TagA", "Hello RocketMQ 1".getBytes(RemotingHelper.DEFAULT_CHARSET)));
messages.add(new Message("TopicTest", "TagA", "Hello RocketMQ 2".getBytes(RemotingHelper.DEFAULT_CHARSET)));
messages.add(new Message("TopicTest", "TagA", "Hello RocketMQ 3".getBytes(RemotingHelper.DEFAULT_CHARSET)));
producer.send(messages);
producer.shutdown();
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroup");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("TopicTest", "*");
consumer.setConsumeMessageBatchMaxSize(32);
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
for (MessageExt message : list) {
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
DefaultMQProducer producer = new DefaultMQProducer("ProducerGroup");
producer.setNamesrvAddr("localhost:9876");
producer.setSendMsgTimeout(10000);
producer.setThreadPoolExecutor(new ThreadPoolExecutor(8, 16, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(2000),
new ThreadFactoryImpl("default")));
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroup");
consumer.setNamesrvAddr("localhost:9876");
consumer.setConsumeThreadMax(64);
consumer.setConsumeThreadMin(20);