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);


上一篇:
下一篇:
切换中文