Day11 RocketMQ集群部署客户端使用

  • 集群部署(由于服务器资源限制2G,选择部署三主)

    我搭建步骤是对的,但是因为太吃内存了,2G支撑不起,所以直接game over

    • 2m-noslave 两主无从节点,存在单点故障【每台机器开一个主节点,勉强能跑】
    • 2m-2s-sync/async 两主两从同步/异步,主节点宕机需手动切换【在每台机器上要开一主一丛,我2G内存受限,会OOM】
    • dledger 具备主从切换的高可用集群。集群中的节点会基于Raft协议随机选举出一个leader【比两主两从还更重,至少需要4G内存】
      • 基于Raft协议,确保多个副本之间特别数据,保证数据强一致性
      • 会接管RocketMQ原生的文件写入,写入速度稍慢,会有性能影响
    bash
    复制代码
    # 将单机部署的机器停机 kill -9 pid # 在机器 A 上 cp /opt/rocketmq/conf/broker.conf /opt/rocketmq/conf/broker-tmp.properties vim /opt/rocketmq/conf/broker-template.properties #brokerClusterName=MyRocketMQCluster #brokerId=0 #deleteWhen=04 #fileReservedTime=48 #brokerRole=ASYNC_MASTER #flushDiskType=ASYNC_FLUSH #listenPort=10911 #storePathRootDir=/opt/rocketmq/store #storePathCommitLog=/opt/rocketmq/store/commitlog # namesrvAddr 先留空,后面统一替换 # 打包mq tar -czf rocketmq.tar.gz rocketmq # 分发到BC两台机器 # 分发到机器 B scp rocketmq.tar.gz root@192.168.1.102:/opt/ # 分发到机器 C scp rocketmq.tar.gz root@192.168.1.103:/opt/ # 解压并删除原文件 cd /opt tar -xzvf rocketmq.tar.gz rm -f rocketmq.tar.gz cd rocketmq mkdir -p store/{commitlog,consumequeue,index,checkpoint,abort} mkdir -p /var/log/rabbitmq # 复制编辑配置文件,三台机器都要改 cp /opt/rocketmq/conf/broker-tmp.properties /opt/rocketmq/conf/broker-c.properties vim /opt/rocketmq/conf/broker-c.properties # brokerName=broker-b # namesrvAddr= IP1:9876;IP2:9876;IP3:9876 # 看内存情况,禁用以下配置,机器4G以上且只部署MQ服务无需操作 # enableScheduleMessageStats=false # timerWheelEnable=false # timerColdDataCheckEnable=false # enableLmq=false # enableMultiDispatch=false # transientStorePoolEnable=false # sendMessageThreadPoolNums=1 # pullMessageThreadPoolNums=1 # 把临时文件删除 rm -f /opt/rocketmq/conf/broker-tmp.properties # 启动nameserver nohup sh /opt/rocketmq/bin/mqnamesrv > /var/log/rocketmq/namesrv.log 2>&1 & # 分别启动broker nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-a.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-b.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-c.properties > /var/log/rocketmq/broker.log 2>&1 & # PS:不出意外的话得出意外了,我2G的内存不够,启动不了,即使禁用了一些配置,AI推荐降级到4.9,但我已经乏了 # 验证集群 # 查看集群 /opt/rocketmq/bin/mqadmin clusterList -n 192.168.1.101:9876 # 应输出: # MyRocketMQCluster broker-a 192.168.1.101:10911 # MyRocketMQCluster broker-b 192.168.1.102:10911 # MyRocketMQCluster broker-c 192.168.1.103:10911 # 测试发送 export NAMESRV_ADDR=192.168.1.101:9876 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer

    image.png

    image.png image.png image.png

  • 运行时架构图(核心:nameServer、broker、client)

    nameServer:服务注册协调中心,独立启动

    broker:提供消息存储、传递、查询等功能,是RocketMQ中最繁琐也是最娇贵的组建

    client:客户端

    💡为什么RocketMQ不使用Zookeeper作为注册中心,而是自己实现nameServer?

    1. 设计简洁且专用 专门的场景设计,更加轻量,易部署
    2. 高性能 zookeeper比较重,繁杂的管理机制和强一致性带来了较多开销
    3. 高可用 无状态,多个nameserver之间对等,即使只有一台示例也不会影响整体运行,而zookpeer在节点之间同步要求严格
    4. 降低依赖性 自己实现可以降低外部系统的依赖,简化系统复杂度,同时能够自己掌控优化
    5. 定制化需求 需要实现特定的功能,不必受限与zookeeper的实现和接口

    image.png

客户端编程模型

  • 消息处理模型

    image.png

  • 消息确认机制

    要支持互联网金融场景,消息安全必须是最高优先级保障。 消息安全有两方面要求,一方面是生产者要能确保将消息发送到Broker上,另一方面消费者要能确保从broker上获取到消息

    • 消息生产端采用消息确认加重试机制保障消息正常发送到RocketMQ

      三种发送消息的方式

      1. 单向发送SendOneWay 只管发送消息,不管成不成功。发送消息效率高,但如果发送失败则无法补救,适用于一些追求效率且允许消息丢失的业务场景。
      2. 同步发送send 发送完消息要等待broker回复结果{SEND_OK,FLUSH_DISK_TIMEOUT,FLUSH_SLAVE_TIMEOUT,SLAVE_NOT_AVAILABLE},保证消息一定发到broker,如果返回失败,可以进行重试,但返回失败不一定代表没有推送给下游消费者。能保障消息发送的安全性,但效率低,会阻塞当前线程。
      3. 异步发送send callback{onSuccess,onException} 不阻塞主线程,由异步回调处理成功或失败的情况,但是在回调结束之前,生产者主线程不能调用shutDown方法关闭主线程。
    • 消费者端采用状态确认机制保证消费者一定能正常处理对应消息

      CONSUME_SUCCESS | RECONSUME_LATER 当broker拿到的是RECONSUME_LATER状态时,会将当前消息放到对应消费者组的重试topic中,以免影响正常队列的运行,当重试次数大于16时,则将消息推入到死信topic,如仍需要重试,可以人工介入进行消息补救

    • 消费者组可以自行指定起始消费位点

      CONSUME_FROM_LAST_OFFSET | CONSUME_FROM_FIRST_OFFSET | CONSUME_FROM_TIMESTAMP 虽然提供了这么一种功能,但是消费者无法知道自己是在那里中断消费的,所以最好是建一个新的消费者组,再从起始点开始消费,也可以指定具有的时间点,consumer.setConsumerTimestamp("20251223171201");

  • 消息处理模式

    • 广播模式 一个消息,推送到所有消费者实例,不关心消费者组
      • consumer.setMessageModel(MessageModel.BROADCASTING);
      • 将offset交由consumer自行保管,只要consumer拉取,就返回对应消息,但是broker端不会对消费失败的消息进行重试
    • 集群模式(默认) 一个消息只会推送到指定的消费者组,由消费者组中的多个实例共同消费
      • 给每个ConsumerGroup维护一个统一的offset,同一个组内只会被消费一次
  • 消息过滤

    • 简单过滤 订阅时使用tag过滤自己感兴趣的内容
      • consumer.subscribe("TagFilterTest", "Tag");
      • 使用两个竖线(||)连接多个Tag值,也可以使用星号(*)匹配所有
    • SQL过滤 使用标准SQL语句进行过滤 WHERE 后面的过滤条件
      • 数值比较,比如:>,>=,<,<=,BETWEEN…AND,=;
      • 字符比较,比如:=,<>,IN;
      • IS NULL 或者 IS NOT NULL;
      • 逻辑符号 AND,OR,NOT;
    • 顺序消息机制
      • 生产者将一组有序的消息发送至同一个MessageQueue,天然保证有序性
      • 消费者实现MessageListenerOrderly接口,按序拿完这一批数据后才去读其他MessageQueue的消息
      • 注意事项
        • 有序性一般是指局部有序,如果需要全局有序,那就只能使用一个queue,性能就很差了
        • 生产者尽可能将有序消息分散到不同的queue,避免数据过于集中导致热点竞争
        • 消费者端只进行有限次数的重试,如果消息处理失败,要保持有序性,就得阻塞一直重试失败消息的消费,但如果一直失败,此时就无法消费其他消息了。
        • 消费者如果处理出现异常,建议返回ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT挂起替代抛出异常
0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP