• Golang 整合RocketMQ


    RocketMQ 相关知识汇总

    RocketMQ 是什么

    RocketMQ 是阿里巴巴开源的一款 MQ 框架,被广泛的使用于不同的业务场景,同时也有非常好的生态系统支持,支持事务消息、顺序消息、批量消息、定时消息、消息回溯等功能。

    RocketMQ核心概念
    • 名称服务(NameServer): 可以理解为注册中心,主要用来保存topic路由消息,管理Broker,在NameServer的集群中,NameServer彼此之间没有任何的通信。

    • 代理服务器(BrokerServer): 消息中转角色,负责存储消息、转发消息。代理服务器在RocketMQ系统中负责接收从生产者发送来的消息并存储、同时为消费者的拉取请求作准备。代理服务器也存储消息相关的元数据,包括消费者组、消费进度偏移和主题和队列消息等。

    • 生产者(Producer):负责生产消息,一般由业务系统负责生产消息。一个消息生产者会把业务应用系统里产生的消息发送到broker服务器。RocketMQ提供多种发送方式,同步发送、异步发送、顺序发送、单向发送。同步和异步方式均需要Broker返回确认信息,单向发送不需要。

    • 生产者组(Producer Group): 同一类Producer的集合,这类Producer发送同一类消息且发送逻辑一致。如果发送的是事务消息且原始生产者在发送之后崩溃,则Broker服务器会联系同一生产者组的其他生产者实例以提交或回溯消费。

    • 消费者(Consumer): 负责消费消息,一般是后台系统负责异步消费。一个消息消费者会从Broker服务器拉取消息、并将其提供给应用程序。从用户应用的角度而言提供了两种消费形式:拉取式消费、推动式消费。

    • 消费者组(Consumer Group): 同一类Consumer的集合,这类Consumer通常消费同一类消息且消费逻辑一致。消费者组使得在消息消费方面,实现负载均衡和容错的目标变得非常容易。要注意的是,消费者组的消费者实例必须订阅完全相同的Topic。RocketMQ支持两种消息模式:集群消费(Clustering)和广播消费(Broadcasting)。

    • 主题(Topic): 表示一类消息的集合,每个主题包含若干条消息,每条消息只能属于一个主题,是RocketMQ进行消息订阅的基本单位。

    • 标签(Tag): 为消息设置的标志,用于同一主题下区分不同类型的消息。来自同一业务单元的消息,可以根据不同业务目的在同一主题下设置不同标签。标签能够有效地保持代码的清晰度和连贯性,并优化RocketMQ提供的查询系统。消费者可以根据Tag实现对不同子主题的不同消费逻辑,实现更好的扩展性。

    扩展概念

    • 消息模型(Message Model): RocketMQ主要由Producer、Broker、Consumer三部分组成,其中Producer负责生产消息,Consumer负责消费消息,Broker负责存储消息。Broker在实际部署过程中对应一台服务器,每个Broker可以存储多个Topic的消息,每个Topic的消息也可以分片存储于不同的Broker。MessageQueue用于存储消息的物理地址,每个Topic中的消息地址存储于多个MessageQueue中。ConsumerGroup由多个Consumer实例构成。

    • 消息(message): 消息系统所传输信息的物理载体,生产和消费数据的最小单位,每条消息必须属于一个主题。RocketMQ中每个消息拥有唯一的Message ID,且可以携带具有业务标识的Key。系统提供了通过Message ID和Key查询消息的功能。

    • 拉取消费(Pull Consumer): Consumer消费的一种类型,应用通常主动调用Consumer的拉消息方法从Broker服务器拉消息、主动权由应用控制。一旦获取了批量消息,应用就会启动消费过程。

    • 推动式消费(Push Consumer): Consumer消费的一种类型,该模式下Broker收到数据后会主动推送给消费端,该消费模式一般实时性较高。

    • 集群消费(Clustering): 集群消费模式下,相同Consumer Group的每个Consumer实例平均分摊消息。

    • 广播消费(Broadcasting): 广播消费模式下,相同Consumer Group的每个Consumer实例都接收全量的消息。

    • 普通顺序消息(Normal Ordered Message): 普通顺序消费模式下,消费者通过同一个消息队列(Topic分区,称作Message Queue)收到的消息是有顺序的,不同消息队列收到的消息则可能是无顺序的。

    • 严格顺序消息(Strictly Ordered Message): 严格顺序消息模式下,消费者收到的所有消息均是有顺序的。

    RocketMQ搭建

    现在我们需要在本地搭建一个rokcetMQ的开发环境,我们搭建的方式是基于docker-compose技术来实现, docker-compose.yaml文件的内容如下:

    version: '3'
    services:
      rmqnamesrv:
        image: rocketmqinc/rocketmq
        container_name: rmqnamesrv
        command: sh mqnamesrv
        ports:
          - "9876:9876"
        volumes:
          - ./namesrv/logs:/root/logs
    
      rmqbroker:
        image: rocketmqinc/rocketmq
        container_name: rmqbroker
        command: sh mqbroker -c /opt/rocketmq-4.9.1/conf/broker.conf
        depends_on:
          - rmqnamesrv
        environment:
          - NAMESRV_ADDR=rmqnamesrv:9876
        ports:
          - "10909:10909"
          - "10911:10911"
        volumes:
          - ./broker/conf:/opt/rocketmq-4.9.1/conf
          - ./broker/logs:/opt/rocketmq-4.9.1/logs
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26

    在创建完成上述内容还需要创建一个rocketmq的配置文件broker.conf文件, 文件的映射路径:./broker/conf,相对于配置文件的路径, 在修改下BrokeIP1的对象地址即可,文件内的内容如下:

    brokerName = broker-a  
    brokerId = 0  
    deleteWhen = 04  
    fileReservedTime = 48  
    brokerRole = ASYNC_MASTER  
    flushDiskType = ASYNC_FLUSH  
    brokerIP1=192.168.18.135
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7

    搭建完成后,使用docker-compose up -d 命令启动就可以了。

    生产者和消费者案例
    • 生产者
    package main
    
    import (
    	"context"
    	"fmt"
    	"github.com/apache/rocketmq-client-go/v2"
    	"github.com/apache/rocketmq-client-go/v2/primitive"
    	"github.com/apache/rocketmq-client-go/v2/producer"
    	"os"
    )
    
    func main() {
    	p, _ := rocketmq.NewProducer(
    		producer.WithNsResolver(primitive.NewPassthroughResolver([]string{"127.0.0.1:9876"})), // 接入地址
    		producer.WithRetry(2),          // 重试次数
    		producer.WithGroupName("test"), // 分组名称
    	)
    	err := p.Start()
    	if err != nil {
    		fmt.Printf("start producer error: %s", err.Error())
    		os.Exit(1)
    	}
    	tags := []string{"TagA", "TagB", "TagC"}
    	for i := 0; i < 3; i++ {
    		tag := tags[i%3]
    		msg := primitive.NewMessage("test",
    			[]byte("Hello RocketMQ Go Client!"))
    		msg.WithTag(tag)
    
    		res, err := p.SendSync(context.Background(), msg)
    		if err != nil {
    			fmt.Printf("send message error: %s\n", err)
    		} else {
    			fmt.Printf("send message success: result=%s\n", res.String())
    		}
    	}
    	err = p.Shutdown()
    	if err != nil {
    		fmt.Printf("shutdown producer error: %s", err.Error())
    	}
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 消费者
    package main
    
    import (
    	"context"
    	"fmt"
    	"github.com/apache/rocketmq-client-go/v2"
    	"github.com/apache/rocketmq-client-go/v2/consumer"
    	"github.com/apache/rocketmq-client-go/v2/primitive"
    )
    
    func main() {
    	c, err := rocketmq.NewPushConsumer(
    		consumer.WithGroupName("test"),
    		consumer.WithNameServer([]string{"127.0.0.1:9876"}),
    	)
    	if err != nil {
    		panic(err)
    	}
    	err = c.Subscribe("test", consumer.MessageSelector{}, func(ctx context.Context,
    		msgs ...*primitive.MessageExt) (consumer.ConsumeResult, error) {
    		for _, msg := range msgs {
    			fmt.Printf("subscribe callback: %+v \n", msg)
    		}
    		return consumer.ConsumeSuccess, nil
    	})
    	if err != nil {
    		panic(err)
    	}
    	err = c.Start()
    	if err != nil {
    		panic(err)
    	}
    	defer func() {
    		err = c.Shutdown()
    		if err != nil {
    
    			fmt.Printf("shutdown Consumer error: %s", err.Error())
    		}
    	}()
    	<-(chan interface{})(nil)
    
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    参考资料
    1. https://zhuanlan.zhihu.com/p/528956421

    2. https://mp.weixin.qq.com/s/iRCP6hEiKOLEp8QRm_OsWQ

  • 相关阅读:
    如何在 java开发工具MyEclipse 中连接到数据库?
    一篇全面而且透彻的RabbitMQ性能优化指南
    IOS企业IPA软件证书 苹果签名证书 有效期到2026年
    外设驱动库开发笔记44:DDC114 ADC驱动
    王道数据结构2(线性表)
    骨灰级大 BOOS 总结出内部不传之密:多线程高并发笔记 + 视频版
    06.论Redis持久化的几种方式
    你不会还不知道B/S与C/S的区别吧?
    java程序连接redis服务器
    【云原生系列第六章】---Serverless架构的应用场景
  • 原文地址:https://blog.csdn.net/weixin_47978762/article/details/134366192