RocketMQ 入门教程:NameServer、Broker、Producer 与 Consumer
· 阅读需 3 分钟
本文讲解前不久进入apache顶级项目的RocketMQ,几个简单的例子讲解如何搭建RocketMQ,以及发送消息,接受消息。
理解 RocketMQ 应先掌握路由发现与消息存储:生产者从 NameServer 获取路由后向 Broker 发送,消费者再按消费组和队列拉取或接收消息。
RocketMQ安装
RocketMQ的安装需要自行编译,接下来编译源码(本文下载源码放在windows系统D:\softwares\目录下)
-
下载源码
git clone -b develop https://github.com/apache/rocketmq.git -
编译
cd rocketmqmvn -Prelease-all -DskipTests clean install -U -
启动rocketmq
cd distribution\target\apache-rocketmqset ROCKETMQ_HOME=D:\softwares\rocketmq\distribution\target\apache-rocketmqbin\mqnamesrv.cmd# 再开启一个cmd窗口,进入到D:\softwares\rocketmq\distribution\target\apache-rocketmqd:cd D:\softwares\rocketmq\distribution\target\apache-rocketmqset ROCKETMQ_HOME=D:\softwares\rocketmq\distribution\target\apache-rocketmqbin\mqbroker.cmd -n localhost:9876mqnamesrv启动成功

mqbroker启动后无输出
PS:不要关闭两个窗口
下面开始代码,写发送消息,接受消息。
引入maven依赖
<dependency><groupId>org.apache.rocketmq</groupId><artifactId>rocketmq-client</artifactId><version>4.1.0-incubating</version></dependency>发送消息
package org.xxz.test.mq;import org.apache.rocketmq.client.exception.MQBrokerException;import org.apache.rocketmq.client.exception.MQClientException;import org.apache.rocketmq.client.producer.DefaultMQProducer;import org.apache.rocketmq.client.producer.SendResult;import org.apache.rocketmq.client.producer.SendStatus;import org.apache.rocketmq.common.message.Message;import org.apache.rocketmq.remoting.exception.RemotingException;public class ProducerTest {public static void main(String[] args) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {DefaultMQProducer producer = new DefaultMQProducer("producerGroup");producer.setNamesrvAddr("127.0.0.1:9876");producer.start();Message msg = new Message();msg.setTopic("test");msg.setBody("hello rocketmq".getBytes());SendResult sendResult = producer.send(msg);if (sendResult.getSendStatus() == SendStatus.SEND_OK) {System.out.println("send msg ok...");}producer.shutdown();}}

接受消息
package org.xxz.test.mq;import java.nio.charset.StandardCharsets;import java.util.List;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;public class ConsumerTest {public static void main(String[] args) throws MQClientException {DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumerGroup");consumer.setNamesrvAddr("127.0.0.1:9876");consumer.subscribe("test", "*");consumer.registerMessageListener(new MessageListenerConcurrently() {@Overridepublic ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {try {for (MessageExt msg : msgs) {System.out.println(new String(msg.getBody(), StandardCharsets.UTF_8));}} catch (Exception e) {return ConsumeConcurrentlyStatus.RECONSUME_LATER;}return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});consumer.start();}}
就是这么简单,后续更精彩哦。。。
相关阅读
需要进一步部署高可用环境时,可以继续阅读 RocketMQ 集群搭建教程;如果关注业务事务一致性,可以阅读 Spring + RocketMQ + MySQL 事务一致性。