ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RocketMQ 入门:Docker 一行起三件套,跑通你的第一条消息

RocketMQ 入门:Docker 一行起三件套,跑通你的第一条消息 作者鱼宵 实战驱动系列 · 第 1 篇完整课程与可运行源码已开源在 Giteehttps://gitee.com/j67mk2/rocketmq-journey 本文对应 lesson-01/下单成功的那一秒系统背后其实还在忙扣积分、发短信、推优惠券……如果这些都让主接口同步去做用户就只能盯着转圈的加载动画干等。这一课我们不堆理论先用 Docker 把 RocketMQ 的 NameServer、Broker、控制台三件套一行命令起好再亲手写两段代码——生产者发一条消息消费者把它收回来。跑通这一条后面的事务消息、延迟消息、顺序消息才有落脚的地方。这些代码我都放在仓库的lesson-01/目录里Gitee 链接见文首clone 下来照着命令一步步重跑一遍比光看印象深十倍——尤其是后面那两个一半正常一半不正常的坑自己亲手踩过一次就忘不掉。一、为什么系统里非要塞个 MQ一句话MQMessage Queue消息队列就是个快递中转站。寄件人把包裹丢给中转站就能走不用自己开车送给收件人收件人什么时候有空什么时候来取。中间多了一站但两边都解放了。它干的就三件事1. 解耦——寄件人和收件人不用互相认识。没有 MQ 时订单系统要通知积分、短信、推荐……每加一个下游订单系统就得改代码多调一次接口。有了 MQ订单系统只管往订单完成这个主题里丢一条消息谁关心谁自己订阅加下游再也不用动订单的代码。2. 异步——点外卖不用站在厨房门口等。用户下单后真正要紧的是订单入库 返回下单成功发短信、发邮件这种慢活让消息在后台慢慢被消费。用户体验从等 2 秒变成等 200 毫秒。3. 削峰——泄洪闸。大促零点的瞬时流量是平时几十倍全砸数据库直接宕机。让请求先涌进 MQ下游按自己的节奏慢慢消费就像水库先蓄水、再匀速放水。MQ 的本质之一用排队换稳定。当然它不是银弹引入 MQ 会带来额外延迟、组件复杂度、消息可能重复消费第 4 课专治可靠性。面试时能补一句什么时候不该用 MQ是加分项。主流 MQ 怎么选先记一张表MQ一句话定位特点RocketMQ国产阿里系事务与可靠性见长事务消息、延迟消息万亿级消息考验中文资料多Kafka日志流事实标准吞吐之王适合日志/指标/流式计算功能偏基础RabbitMQ轻量灵活路由AMQP 协议、路由灵活社区老牌吞吐相对低选型一句话日志流找 Kafka轻量灵活找 RabbitMQ涉及事务/可靠性的业务消息优先 RocketMQ——这也是本系列选它的原因。二、四个组件 Topic/Queue一张图先建立全局先看图别慌每个词下面都有人话解释┌─────────────────────────────┐ │ NameServer电话簿 │ │ 无状态注册中心Broker 在哪 │ └──────────────┬──────────────┘ │ ① 查询 Broker 地址 ┌──────────┐ │ ┌──────────┐ │ Producer │─发消息─────▼──存储转发──▶│ Consumer │ │ (寄件人) │ ┌──────────┐ │ (收件人) │ └──────────┘ │ Broker │ └──────────┘ │ (邮局仓库) │ │ Topic Queue0 消息1,消息2 │ Queue1 消息3 └──────────┘逐个对号入座组件类比一句话职责Producer生产者寄件人发消息的人把消息发到 Broker 的某个 TopicConsumer消费者收件人收消息的人从队列里拉取消息处理Broker代理邮局仓库真正存消息、转发消息的服务消息落盘在它的磁盘上NameServer命名服务电话簿无状态注册中心记录哪个 Broker 在哪给客户端指路它自己不存消息Topic主题邮筒分类消息的类别比如订单消息物流消息各一个 TopicQueue队列货架通道Topic 下面真正存消息的物理队列一个 Topic 默认 4 个队列队列是并行度的最小单位ConsumerGroup消费组同一班快递员一组消费者一条消息只会被组内一个消费者消费第 3 课细讲这里最容易被绕晕的是 Topic 和 Queue 的关系Topic 是逻辑分类你跟别人说这是订单消息Queue 才是物理上真正存消息的地方。一个 Topic 默认拆成 4 个 Queue 并行跑消息一条一条轮流落进去——这既是并行度也是后面顺序消息为什么难的根源第 5 课讲。一条消息的完整旅程就五步Producer 先问 NameServer“TopicLesson01 在哪个 Broker”NameServer 回答“在 broker-a地址是 xxx:10911”Producer 连上 Broker把消息写进 Topic 的某个 QueueConsumer 同样问 NameServer 拿到地址订阅这个 TopicBroker 把 Queue 里的消息推给消费者处理记住一句话NameServer 挂了只是暂时找不到路Broker 挂了消息才真的危险第 9 课讲 Broker 高可用。三、动手Docker 一行起三件套环境Windows Docker会docker ps就行。RocketMQ 的 Broker 对 JDK 版本敏感用官方镜像包好环境一条命令起全套不污染系统。第 1 步30 秒检查环境。java-version# 期望看到 17.0.xdocker version# Client 和 Server 都有版本号docker compose version# 期望 v2 及以上第 2 步进入课程根目录一条命令起三件套。cd 你的课程根目录\rocketmq-journey# 也就是 docker-compose.yml 所在目录docker compose up-d这条命令读根目录的docker-compose.yml依次拉起三个容器rmqnamesrvNameServer对外端口9876rmqbrokerBroker对外端口10911rmqdashboard控制台浏览器访问http://localhost:8080首次会拉镜像几百 MB等几分钟以后启动就是秒级。第 3 步确认真的起来了。docker composeps# 三行容器都是 Up 状态再确认 Broker 已经登记到电话簿里docker exec rmqbroker sh mqadmin clusterList-n namesrv:9876看到一行类似DefaultCluster broker-a 0 172.26.209.90:10911 V4_9_7 ...就说明 Broker 已上线。最后浏览器打开http://localhost:8080左侧能看到集群/主题/消费者菜单三件套全部就绪。Topic 怎么建两种方式先有个印象本课用的是自动创建——我们在 Broker 配置里开了autoCreateTopicEnabletrue生产者第一次往TopicLesson01发消息时它自动建好默认 4 个队列适合学习开发。生产环境一般关掉自动创建改用控制台或mqadmin命令显式建防止手滑写错 Topic 名建出一堆垃圾主题。四、写代码发一条、收一条工程是个纯 Maven 项目lesson-01 目录。先编译cd 你的课程根目录\rocketmq-journey\lesson-01$env:JAVA_HOMEC:\jdk\jdk-17.0.12# Maven 必须跑在 JDK 17 上每个新窗口设一次mvn clean install-DskipTests# 结尾看到 BUILD SUCCESS 即可先看生产者 SyncProducer——干的事就是往TopicLesson01同步发 2 条消息packagecom.example.lesson01;importorg.apache.rocketmq.client.producer.DefaultMQProducer;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.common.message.Message;importjava.nio.charset.StandardCharsets;/** * 第 1 课 Demo同步发送消息的生产者 * 干的事向 TopicLesson01 这个主题发送 2 条普通消息。 * 用同步发送发一条、等 Broker 确认一条能立刻知道成没成功。 */publicclassSyncProducer{publicstaticvoidmain(String[]args)throwsException{// 1. 创建生产者。构造参数是生产者组名同组生产者共享同一套发送配置DefaultMQProducerproducernewDefaultMQProducer(lesson01_producer_group);// 2. 告诉生产者 NameServer 在哪。// NameServer 相当于电话簿客户端先问它 Broker 在哪再去连 Broker。// 本地 Docker 起的 namesrv 容器对外端口是 9876producer.setNamesrvAddr(127.0.0.1:9876);// 3. 启动生产者内部会去 NameServer 拉取 Broker 地址并建立连接producer.start();// 4. 循环发送 2 条消息for(inti0;i2;i){// 构造消息三要素 Topic发到哪个主题 Tag给消息打的小标签用于过滤 消息体StringbodyHello RocketMQ这是第 (i1) 条消息;MessagemsgnewMessage(TopicLesson01,TagA,body.getBytes(StandardCharsets.UTF_8));// 5. 同步发送一直等到 Broker 确认收到并落盘返回 SendResultSendResultresultproducer.send(msg);// SendResult 里能看到消息落在哪个队列QueueId、队列里的偏移量QueueOffsetSystem.out.println(发送成功结果result);}// 6. 用完关闭生产者释放连接producer.shutdown();System.out.println(生产者已关闭。);}}记住三个第一反应要写的东西组名lesson01_producer_group、NameServer 地址127.0.0.1:9876、TopicTag消息体。producer.send()是同步发送发完必须等 Broker 回执sendStatusSEND_OK才算真成功——这是最稳的发送姿势第 2 课对比同步、异步、单向。跑生产者mvn exec:java-Dexec.mainClasscom.example.lesson01.SyncProducer再看消费者 PushConsumer——订阅TopicLesson01里 TagA 的消息收到就打印packagecom.example.lesson01;importorg.apache.rocketmq.client.consumer.DefaultMQPushConsumer;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;importorg.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;importorg.apache.rocketmq.common.consumer.ConsumeFromWhere;importorg.apache.rocketmq.common.message.MessageExt;importjava.nio.charset.StandardCharsets;importjava.util.List;importjava.util.concurrent.atomic.AtomicInteger;/** * 第 1 课 DemoPush 模式消费者 * Push模式 消息由 Broker 主动推给消费者消费者只管注册一个回调。 * 教程演示收到 2 条后或最多等 25 秒自动退出生产中消费者是常驻进程。 */publicclassPushConsumer{publicstaticvoidmain(String[]args)throwsException{// 1. 创建消费者。参数是消费组名同组内一条消息只会被一个消费者消费第 3 课细讲DefaultMQPushConsumerconsumernewDefaultMQPushConsumer(lesson01_consumer_group);// 2. 同样要告诉消费者 NameServer 在哪consumer.setNamesrvAddr(127.0.0.1:9876);// 3. 从最早的消息开始消费这样先发消息、后启动消费者也能把历史消息捞回来consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);// 4. 订阅只关心 TopicLesson01 里 tag 为 TagA 的消息写 * 表示收全部consumer.subscribe(TopicLesson01,TagA);// 5. 注册消息监听器回调Broker 把消息推过来时这个回调会被调用AtomicIntegerreceivedCountnewAtomicInteger(0);consumer.registerMessageListener(newMessageListenerConcurrently(){OverridepublicConsumeConcurrentlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){StringbodyTextnewString(msg.getBody(),StandardCharsets.UTF_8);System.out.println(收到消息: bodyTexttagmsg.getTags());receivedCount.incrementAndGet();}// 返回 CONSUME_SUCCESS 表示处理成功可以确认返回 RECONSUME_LATER 会触发重试returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});// 6. 启动消费者内部会去 NameServer 找 Broker并开始拉取消息consumer.start();System.out.println(消费者已启动正在等待消息……收到 2 条后自动退出最多等 25 秒);// 7. 演示用退出逻辑收到 2 条或 25 秒超时就关闭消费者longdeadlineSystem.currentTimeMillis()25_000;while(receivedCount.get()2System.currentTimeMillis()deadline){Thread.sleep(500);}consumer.shutdown();System.out.println(观察结束消费者已关闭。);}}两个伏笔先埋着CONSUME_FROM_FIRST_OFFSET是从最早消息开始消费因为我们是先发消息、后开消费者没这句就收不到历史消息CONSUME_SUCCESS / RECONSUME_LATER决定了消息是确认还是重试第 3 课细讲。跑消费者mvn exec:java-Dexec.mainClasscom.example.lesson01.PushConsumer五、实测输出SendResult 到底说了什么下面是同款环境Windows 11 Docker 29 JDK 17的真实运行输出。先跑生产者发送成功结果SendResult [sendStatusSEND_OK, msgId7F000001CB6036CC9385090CB0C40000, offsetMsgIdAC1AD15A00002A9F0000000000000000, messageQueueMessageQueue [topicTopicLesson01, brokerNamebroker-a, queueId0], queueOffset0] 发送成功结果SendResult [sendStatusSEND_OK, msgId7F000001CB6036CC9385090CB0F30001, offsetMsgIdAC1AD15A00002A9F00000000000000D8, messageQueueMessageQueue [topicTopicLesson01, brokerNamebroker-a, queueId1], queueOffset0] 生产者已关闭。逐条看懂 SendResult注意两条消息落在了 queueId0 和 queueId1这是正常的——消息会自动散到不同队列字段含义sendStatusSEND_OK发送成功还有刷盘超时、主从超时、从节点不可用三种失败态第 4 课讲msgId客户端生成的全局唯一 ID可用于追踪这条消息offsetMsgIdBroker 落盘后生成的物理偏移 IDqueueId / queueOffset消息落在哪个 Topic 的第几号队列、第几个位置再跑消费者消费者已启动正在等待消息……收到 2 条后自动退出最多等 25 秒 收到消息: Hello RocketMQ这是第 1 条消息tagTagA 收到消息: Hello RocketMQ这是第 2 条消息tagTagA 观察结束消费者已关闭。最后去浏览器http://localhost:8080 → 主题搜TopicLesson01能看到这个 Topic 下面 4 个队列QueueId 0~3和消息总量。到这里你人生中第一条 RocketMQ 消息就跑通了。排查提示客户端的详细日志不打印在控制台而是写到C:\Users\你\logs\rocketmqlogs\rocketmq_client.log连了哪个 Broker、心跳、拉取全过程都在里面。六、新手最容易卡住的两个一半正常一半不正常第一课最劝退的不是概念而是这两个诡异现象提前打个预防针。坑 1Broker 注册了客户端却连不上——brokerIP1 得填对地址。NameServer 里能看到 Broker但 Java 客户端发消息报connect to ... failed。原因是 Windows Docker DesktopWSL2下端口只绑在回环地址和 WSL 虚拟机 IP 上Broker 注册给 NameServer 的地址必须是客户端真能连上的那个。把broker.conf里的brokerIP1改成wsl hostname -I查到的虚拟机 IP再docker compose restart broker即可。电脑重启后 WSL 的 IP 会变变了就重复一遍——这就是很多人说Windows 上跑 RocketMQ 麻烦的真相。坑 2生产者发得好好的消费者一条都收不到——JDK 在 cgroup 下崩了。发送是SEND_OK控制台也正常消费者却一直收不到客户端日志刷NoClassDefFoundError: Could not initialize class ...StoreUtil。根因是 Broker 处理拉取请求时要读物理内存镜像里的 JDK 8 在 WSL2 的 cgroup v2 环境下初始化内存探测直接抛异常把StoreUtil这个类搞崩了而发送路径不碰它所以一半正常一半不正常。解决办法是给 Broker 的 JVM 加-XX:-UseContainerSupport禁用 cgroup 探测然后docker compose up -d重建容器。排查口诀发送成功但消费不到先去看 Broker 服务端的 remoting.log别在客户端反复折腾。七、挑战题答案在仓库源码里跑起来才知道⭐ 把SyncProducer改成循环发 10 条消息里带上循环下标再跑PushConsumer观察 10 条是否全收到、收到顺序和发送顺序是否一致。⭐⭐ 在控制台http://localhost:8080手动新建一个TopicHomework014 个队列改代码把消息发到这个新 Topic验证控制台能看到新 Topic 的队列和消息量。⭐⭐⭐ 先不启动消费者跑生产者发 10 条等 1 分钟后把消费者里的CONSUME_FROM_FIRST_OFFSET去掉、用默认值再跑观察还收不收得到——这背后是消费进度offset从哪开始的问题第 3 课彻底讲。八、面试回答模板面试官为什么系统要引入 MQ一句话解耦、异步、削峰本质是用排队换稳定。解耦是下游变化不改上游代码异步是把慢操作挪到后台、缩短主接口响应削峰是用 MQ 扛瞬时洪峰、下游匀速消费。代价是引入额外延迟、复杂度和重复消费问题不是什么场景都该上。见本文第一节追问RocketMQ 有哪些核心组件Producer 发消息、Consumer 收消息、Broker 存消息转发消息、NameServer 是无状态注册中心给客户端指路。Topic 是逻辑分类Queue 是 Topic 下真正存消息的物理队列一个 Topic 默认 4 个队列。见本文第二节追问NameServer 和 Broker 有什么区别NameServer 是电话簿、无状态、不存消息挂了只是暂时找不到路Broker 才是消息仓库消息落盘在它磁盘上它挂了消息才真危险。追问同步发送怎么判断成功看返回的SendResult.sendStatus是不是SEND_OK而不是没抛异常。关于这个系列本文是「Java 后端实战精通营」系列第 1 篇原则实战驱动、由浅到深、面试向每篇文章的结论都可以亲手验证。RocketMQ 实战精通营10 课https://gitee.com/j67mk2/rocketmq-journey本文对应源码位置lesson-01/最小 Maven 工程内含SyncProducer同步生产者 PushConsumer推送消费者系列文章一览按发布顺序篇主题1RocketMQ 入门Docker 一行起三件套跑通你的第一条消息2RocketMQ 发送方式同步异步批量单向消息都怎么发出去3RocketMQ 消费模式集群、广播与重试消息怎么被吃掉4RocketMQ 可靠性发送重试加幂等消息一条都不丢5RocketMQ 顺序消息订单流程不乱套的秘密6RocketMQ 事务消息订单与积分的最终一致7RocketMQ 延迟消息30 分钟未支付自动关单怎么做8RocketMQ 积压治理百万消息堵在队列怎么办9RocketMQ 集群高可用与过滤主从架构 Tag 精准投递10RocketMQ 面试冲刺高频考点一口气背完下一篇预告《RocketMQ 发送方式同步异步批量单向消息都怎么发出去》——本课只用了最稳的同步发送下一课把同步、异步、单向三种姿势讲全再带你发一条10 秒后才被消费的延迟消息。跑完有任何报错把终端输出发评论区一起排查。
返回列表