多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

RocketMQ 集群高可用与过滤:主从架构 + Tag 精准投递

RocketMQ 集群高可用与过滤:主从架构 + Tag 精准投递 RocketMQ 集群高可用与过滤主从架构 Tag 精准投递作者鱼宵 实战驱动系列 · 第 9 篇完整课程与可运行源码已开源在 Giteehttps://gitee.com/j67mk2/rocketmq-journey 本文对应 lesson-09/一、两个真实场景把这篇文章的问题先抛给你场景一订单 Topic 里混着三种消息——支付成功、退款通知、发货提醒。下游的积分系统只该给支付成功的用户加积分结果它把退款通知也消费了用户明明退款了积分反而涨了。怎么让 Broker 只把支付成功这一类消息推给它别的根本别送过来场景二你把这套消息队列部署成一台 NameServer 一台 Broker学习时好好的上了生产这台 Broker 凌晨挂了——消息发不进来、消费者取不到整个下单链路直接停摆。更坑的是有同学照着网上教程写 SQL92 过滤消费者刚 start() 就抛The broker does not support consumer to filter message by SQL92对着报错一脸懵。这两件事一个是消息怎么精准投递一个是集群怎么不挂。这篇文章两件事一起讲高可用部分主从、Dledger以概念 类比为主代码部分我们动手把Tag 过滤完整跑一遍并把那个 SQL92 报错的真实原因挖出来。答案都在 lesson-09 的源码里clone 下来跑一遍就知道。二、集群高可用为什么不能只有一台机器前面 8 课我们都是单 NameServer 单 Broker。这两个角色挂了的代价完全不一样NameServer 挂了它是电话簿只记着 Broker 在哪自己不存消息。它无状态客户端配置多个地址挂一个自动换另一个消息不丢。Broker 挂了消息存在它的磁盘上这才是要命的——发不进、取不到业务停摆。NameServer 多实例多印几本电话簿NameServer 无状态、实例之间互不通信所以多起几个几乎零成本NameServer-1 (电话簿A) NameServer-2 (电话簿B) \ / \ / producer.setNamesrvAddr(ns1:9876;ns2:9876)客户端配置里写两个地址分号隔开连不上这个就连那个。一句话NameServer 高可用 多起几个客户端配多个地址它挂了不丢消息只影响找路。Broker 主从复制一个正仓库一个副仓库Broker 的高可用靠主从复制。类比一个仓库有正仓库MasterbrokerId0和副仓库SlavebrokerId1写消息 ──▶ Master(正仓库) ──复制──▶ Slave(副仓库) 负责读写 只备份、接读请求Master 对外提供读写消息先写它Slave 实时从 Master 复制一份自己不接写请求。主从之间怎么同步两种模式面试高频模式怎么复制优点缺点类比SYNC_MASTER 同步复制Master 写后等 Slave 也写成功才返回生产者发送成功消息不丢至少两处有副本慢要等副仓库确认寄快递必须正副仓库都签收才告诉你发货成功ASYNC_MASTER 异步复制Master 写完自己就返回成功后台再异步同步给 Slave快不等 SlaveMaster 刚写完就挂、还没同步的那一小撮消息可能丢正仓库先签收说发走了回头再慢慢抄给副仓库怎么选一句话金融、订单这种不能丢的业务用同步复制容忍极小概率丢失、追求吞吐的用异步复制。部署形态与 Dledger主挂了谁来当新主普通主从有个痛点Master 挂了Slave 虽然有完整数据但没人加冕它成新 Master整个组就停写了得人工介入。Dledger 就是来解决这个的——它基于 Raft 协议一种分布式一致性算法让一组通常 3 个节点自己投票选一个 Leader 负责写Leader 挂了剩下的自动再投票选新 Leader几十秒内自动恢复写服务全程不用人管。形态结构说明单主1 个 Broker学习测试用挂了全停本课 docker-compose 就是这种一主一从1 主 1 从主挂了从能顶上读但普通模式不能自动切主写双主双从2 个 Master 各带 1 个 Slave生产常见两个主各管一部分队列互为备份Dledger 多副本3 节点成组Raft 自动选主主挂了自动选新主故障自动转移面试记一句Dledger 用 Raft 把主从复制 人工切主升级成多副本 自动选主。三、消息过滤收件人只想要某几种快递高可用讲完部署部分本课只讲不跑进入能动手的部分。一个 Topic 往往混着多种消息不同下游只关心其中一部分。消息过滤就是让消费者只收到自己关心的那部分——类比你给收件地址备注只收顺丰和京东其他放门卫快递员Broker送之前先按备注筛一遍。RocketMQ 提供两种过滤方式方式一Tag 过滤最常用。发消息时打 Tag消费时订阅只写要哪个consumer.subscribe(TopicLesson09,TagA);// 只要 TagAconsumer.subscribe(TopicLesson09,TagA || TagB);// TagA 或 TagB 都要consumer.subscribe(TopicLesson09,*);// 全收方式二SQL92 属性过滤更灵活。发消息时用putUserProperty塞任意用户属性消费时写类 SQL 条件msg.putUserProperty(age,20);msg.putUserProperty(city,beijing);consumer.subscribe(TopicLesson09,MessageSelector.bySql(age 18 AND city beijing));两者对比对比Tag 过滤SQL92 属性过滤过滤依据只有 Tag 一个字段任意用户属性age、city……表达能力等于、或TagA||TagB比较大小、and/or、like接近 where 子句性能最好Broker 端简单匹配稍差要执行表达式限制几乎都能用需要 Broker 开 enablePropertyFiltertrue四、动手第一步发 6 条带 Tag 和属性的消息先把服务起好docker compose up -d然后编译本课工程。生产者交替发 TagA系统通知和 TagB营销推送各 3 条每条都带上 age/city 两个属性packagecom.example.lesson09;importorg.apache.rocketmq.client.producer.DefaultMQProducer;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.common.message.Message;importjava.nio.charset.StandardCharsets;/** * 第 9 课 Demo发送带 Tag 和用户属性的消息过滤的原料 */publicclassTagFilterProducer{publicstaticvoidmain(String[]args)throwsException{DefaultMQProducerproducernewDefaultMQProducer(lesson09_producer_group);producer.setNamesrvAddr(127.0.0.1:9876);producer.start();// 交替发 6 条TagA系统通知/ TagB营销推送各 3 条String[]tags{TagA,TagB,TagA,TagB,TagA,TagB};for(inti0;itags.length;i){Stringtagtags[i];Stringbody(tag.equals(TagA)?系统通知-:营销推送-)(i1);MessagemsgnewMessage(TopicLesson09,tag,body.getBytes(StandardCharsets.UTF_8));// 塞用户属性SQL92 过滤就靠它。奇数 age20/beijing偶数 age16/shanghaiintage(i%20)?20:16;Stringcity(i%20)?beijing:shanghai;msg.putUserProperty(age,String.valueOf(age));msg.putUserProperty(city,city);SendResultresultproducer.send(msg);System.out.println(发送成功: [tag] bodyageage, citycityqueueIdresult.getMessageQueue().getQueueId());}producer.shutdown();System.out.println(发送完成共 6 条TagA/TagB 各 3 条带 age/city 属性。);}}两个关键点new Message(topic, tag, body)这一步把 Tag 打上Tag 过滤能不能用全看它putUserProperty塞的属性是 SQL92 过滤的原料——你不在消息上塞属性后面bySql(age 18)就没东西可比。五、动手第二步只订阅 TagA实测看过滤效果消费者只订阅 TagA关键就是subscribe的第二个参数packagecom.example.lesson09;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;/** * 第 9 课 DemoTag 过滤消费者只收 TagA。 */publicclassTagFilterConsumer{publicstaticvoidmain(String[]args)throwsException{DefaultMQPushConsumerconsumernewDefaultMQPushConsumer(lesson09_consumer_group);consumer.setNamesrvAddr(127.0.0.1:9876);consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);// 关键第二个参数 TagA 只收 TagA。* 全收TagA || TagB 两个都要consumer.subscribe(TopicLesson09,TagA);AtomicIntegerreceivedCountnewAtomicInteger(0);consumer.registerMessageListener(newMessageListenerConcurrently(){OverridepublicConsumeConcurrentlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){StringbodyTextnewString(msg.getBody(),StandardCharsets.UTF_8);// 把 tag 和属性都打出来方便核对是不是真的只有 TagASystem.out.println(收到: bodyTexttagmsg.getTags(), agemsg.getUserProperty(age), citymsg.getUserProperty(city));receivedCount.incrementAndGet();}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});consumer.start();System.out.println(Tag 过滤消费者已启动只订阅 TagA……收满 3 条或 25 秒后自动退出);longdeadlineSystem.currentTimeMillis()25_000;while(receivedCount.get()3System.currentTimeMillis()deadline){Thread.sleep(500);}consumer.shutdown();System.out.println(观察结束共收到 receivedCount.get() 条应该全是 TagA消费者已关闭。);}}过滤就发生在subscribe(TopicLesson09, TagA)这一行的第二个参数——Broker 端按这个表达式筛选不匹配的消息根本不推给你。下面是本机真实运行输出先跑生产者再跑消费者生产者发消息的输出发送成功: [TagA] 系统通知-1age20, citybeijingqueueId2 发送成功: [TagB] 营销推送-2age16, cityshanghaiqueueId3 发送成功: [TagA] 系统通知-3age20, citybeijingqueueId0 发送成功: [TagB] 营销推送-4age16, cityshanghaiqueueId1 发送成功: [TagA] 系统通知-5age20, citybeijingqueueId2 发送成功: [TagB] 营销推送-6age16, cityshanghaiqueueId3 发送完成共 6 条TagA/TagB 各 3 条带 age/city 属性。消费者只订阅 TagA 的输出关键证据Tag 过滤消费者已启动只订阅 TagA……收满 3 条或 25 秒后自动退出 收到: 系统通知-3tagTagA, age20, citybeijing 收到: 系统通知-1tagTagA, age20, citybeijing 收到: 系统通知-5tagTagA, age20, citybeijing 观察结束共收到 3 条应该全是 TagA没有 TagB消费者已关闭。读懂这次输出现象说明收到的全是系统通知TagATag 过滤生效TagB营销推送一条都没推过来顺序是 3、1、5 而不是 1、3、5消息打散在多个队列、多消费线程并发处理顺序不保证第 8 课讲过并发消费age/city 属性也带出来了消费时能读回生产者塞的属性为 SQL92 过滤做铺垫六、SQL92 属性过滤代码看懂但先看那个真实报错SQL92 过滤的写法是把 subscribe 的第二个参数换成MessageSelector.bySql(...)packagecom.example.lesson09;importorg.apache.rocketmq.client.consumer.DefaultMQPushConsumer;importorg.apache.rocketmq.client.consumer.MessageSelector;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;/** * 第 9 课 DemoSQL92 属性过滤消费者。条件age 18 AND city beijing。 */publicclassSqlFilterConsumer{publicstaticvoidmain(String[]args)throwsException{// 用全新消费组避免和 TagFilterConsumer 的 offset 互相影响DefaultMQPushConsumerconsumernewDefaultMQPushConsumer(lesson09_sql_consumer_group);consumer.setNamesrvAddr(127.0.0.1:9876);consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);// 关键用 MessageSelector.bySql(...) 代替 tag 参数写类 SQL 条件consumer.subscribe(TopicLesson09,MessageSelector.bySql(age 18 AND city beijing));AtomicIntegerreceivedCountnewAtomicInteger(0);consumer.registerMessageListener(newMessageListenerConcurrently(){OverridepublicConsumeConcurrentlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){StringbodyTextnewString(msg.getBody(),StandardCharsets.UTF_8);System.out.println(收到(SQL过滤): bodyTexttagmsg.getTags(), agemsg.getUserProperty(age), citymsg.getUserProperty(city));receivedCount.incrementAndGet();}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});consumer.start();System.out.println(SQL92 过滤消费者已启动条件age 18 AND city beijing……最多观察 20 秒自动退出);longdeadlineSystem.currentTimeMillis()20_000;while(System.currentTimeMillis()deadline){Thread.sleep(500);}consumer.shutdown();System.out.println(观察结束共收到 receivedCount.get() 条满足条件的消息消费者已关闭。);}}但在本课的单机教学 Broker 上这个程序启动就报错真实输出org.apache.rocketmq.client.exception.MQClientException: CODE: 1 DESC: The broker does not support consumer to filter message by SQL92 at com.example.lesson09.SqlFilterConsumer.main (SqlFilterConsumer.java:66)原因一句话SQL92 的 Broker 端过滤不是默认开的需要在 broker.conf 里加一行enablePropertyFiltertrue重建 Broker 再试。教学环境只跑通最常用的 Tag 过滤所以这里看懂代码写法即可——真把开关打开满足age18 AND citybeijing的恰好就是系统通知 1/3/5 那三条age20、citybeijing。顺手记四个这课真实踩过的小坑坑原因解法SQL92 启动就报 does not support SQL92Broker 没开 enablePropertyFilterbroker.conf 加一行开关再重建import 写成 …consumer.selector.MessageSelector 编译就挂凭印象记错包名4.9.7 里它在 org.apache.rocketmq.client.consumer.MessageSelectorimport 写对包名记不住去 jar 里搜类名两个消费者用同一消费组offset 互相打架同组共享一份 offset书签上一组消费到末尾新消费者从书签后等新消息不同实验各用各的消费组名Tag 多写个空格就一条都收不到Tag 表达式是精确匹配TagA ≠ “TagA”别手抖加空格多个 Tag 用 || 连接别写逗号七、挑战题答案在仓库源码里跑起来才知道⭐ 把TagFilterConsumer的订阅从TagA改成TagA || TagB重新发消息再消费看看是不是 6 条全收到再改成*对比三者区别。⭐⭐ 给生产者发的消息多加一个属性levelvip/normal写一个 SQL92 条件level vip AND age 18。思考要让它真跑起来需要改 Broker 的哪个配置提示第六节那个报错⭐⭐⭐ 画一张双主双从 3 节点 NameServer部署图标清楚生产者连哪个地址、Master 和 Slave 谁负责写、一个 Master 挂了读和写分别会发生什么。这张图能口头讲清楚高可用部分就算真学会了。八、面试回答模板面试官RocketMQ 怎么保证高可用三件套NameServer 多实例无状态挂了换另一本电话簿不丢消息 Broker 主从复制正仓库写、副仓库备份 Dledger基于 Raft 的多副本组主挂了自动投票选新主。展开见本文第二节。追问同步复制和异步复制怎么选同步SYNC_MASTER主写完等从也写好才返回稳、不丢、慢异步ASYNC_MASTER主写完就返回、后台再同步快、但主刚写就挂可能丢一小段。金融订单用同步追求吞吐用异步。追问Tag 过滤和 SQL92 过滤有什么区别SQL92 跑不起来查什么Tag 按标签做等于/或运算简单性能好几乎都能用SQL92 按用户属性写 where 条件灵活但需要 Broker 开 enablePropertyFiltertrue。跑不起来依次查Broker 开关开了没、消息有没有 putUserProperty、SQL 语法对不对。追问普通主从 Master 挂了会怎样Slave 有完整数据但普通模式不能自动当主写会停写需要人工或 Dledger 介入切主——这正是 Dledger 存在的意义。九、关于这个系列本文是「Java 后端实战精通营」系列第 9 篇原则实战驱动、由浅到深、面试向每篇文章的结论都可以亲手验证。RocketMQ 实战精通营10 课https://gitee.com/j67mk2/rocketmq-journey本文对应源码位置lesson-09/三个可运行主类TagFilterProducer 发消息、TagFilterConsumer 验 Tag 过滤、SqlFilterConsumer 演示 SQL92 写法系列文章一览按发布顺序篇主题1RocketMQ 入门Docker 一行起三件套跑通你的第一条消息2RocketMQ 发送方式同步异步批量单向消息都怎么发出去3RocketMQ 消费模式集群、广播与重试消息怎么被吃掉4RocketMQ 可靠性发送重试加幂等消息一条都不丢5RocketMQ 顺序消息订单流程不乱套的秘密6RocketMQ 事务消息订单与积分的最终一致7RocketMQ 延迟消息30 分钟未支付自动关单怎么做8RocketMQ 积压治理百万消息堵在队列怎么办9RocketMQ 集群高可用与过滤主从架构 Tag 精准投递10RocketMQ 面试冲刺高频考点一口气背完下一篇预告《RocketMQ 面试冲刺高频考点一口气背完》——把 9 课高频面试题整理成 30 问速记清单配上简历包装话术帮你把这趟学习变成面试桌上的筹码。跑完有任何报错把终端输出发评论区一起排查。
返回列表