<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>3.6.2</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> <version>2.5.4</version> </dependency> rocketmq.name-server=127.0.0.1:9876 rocketmq.producer.group=mygroup rocketmq.consumer.group=mygroup @Configuration public class RocketMQConfig { @Value("${rocketmq.name-server}") private String nameServer; @Value("${rocketmq.producer.group}") private String producerGroup; @Bean public DefaultMQProducer rocketMQProducer() { DefaultMQProducer producer = new DefaultMQProducer(producerGroup); producer.setNamesrvAddr(nameServer); return producer; } } @Service public class RocketMQService { @Autowired private DefaultMQProducer producer; public void sendMessage(String topic, String message) throws Exception { Message msg = new Message(topic, "TagA", message.getBytes(StandardCharsets.UTF_8)); SendResult sendResult = producer.send(msg); System.out.println(sendResult); } } @Configuration public class RocketMQConfig { @Value("${rocketmq.name-server}") private String nameServer; @Value("${rocketmq.consumer.group}") private String consumerGroup; @Bean public DefaultMQPushConsumer rocketMQConsumer() throws MQClientException { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup); consumer.setNamesrvAddr(nameServer); consumer.subscribe("topic", "*"); consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) { for (MessageExt message : messages) { System.out.println(new String(message.getBody())); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); return consumer; } }


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