亚洲乱码中文字幕综合,中国熟女仑乱hd,亚洲精品乱拍国产一区二区三区,一本大道卡一卡二卡三乱码全集资源,又粗又黄又硬又爽的免费视频

go-micro集成RabbitMQ實(shí)戰(zhàn)和原理詳解

 更新時(shí)間:2022年05月07日 08:53:13   作者:波斯馬  
本文主要介紹go-micro使用RabbitMQ收發(fā)數(shù)據(jù)的方法和原理,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下

在go-micro中異步消息的收發(fā)是通過(guò)Broker這個(gè)組件來(lái)完成的,底層實(shí)現(xiàn)有RabbitMQ、Kafka、Redis等等很多種方式,這篇文章主要介紹go-micro使用RabbitMQ收發(fā)數(shù)據(jù)的方法和原理。

Broker的核心功能

Broker的核心功能是Publish和Subscribe,也就是發(fā)布和訂閱。它們的定義是:

Publish(topic string, m *Message, opts ...PublishOption) error
Subscribe(topic string, h Handler, opts ...SubscribeOption) (Subscriber, error)

發(fā)布

發(fā)布第一個(gè)參數(shù)是topic(主題),用于標(biāo)識(shí)某類消息。

發(fā)布的數(shù)據(jù)是通過(guò)Message承載的,其包括消息頭和消息體,定義如下:

type Message struct {
	Header map[string]string
	Body   []byte
}

消息頭是map,也就是一組KV(鍵值對(duì))。

消息體是字節(jié)數(shù)組,在發(fā)送和接收時(shí)需要開(kāi)發(fā)者進(jìn)行編碼和解碼的處理。

訂閱

訂閱的第一個(gè)參數(shù)也是topic(主題),用于過(guò)濾出要接收的消息。

訂閱的數(shù)據(jù)是通過(guò)Handler處理的,Handler是一個(gè)函數(shù),其定義如下:

type Handler func(Event) error

其中的參數(shù)Event是一個(gè)接口,需要具體的Broker來(lái)實(shí)現(xiàn),其定義如下:

type Event interface {
	Topic() string
	Message() *Message
	Ack() error
	Error() error
}
  • Topic() 用于獲取當(dāng)前消息的topic,也是發(fā)布者發(fā)送時(shí)的topic。
  • Message() 用于獲取消息體,也是發(fā)布者發(fā)送時(shí)的Message,其中包括Header和Body。
  • Ack() 用于通知Broker消息已經(jīng)收到了,Broker可以刪除消息了,可用來(lái)保證消息至少被消費(fèi)一次。
  • Error() 用于獲取Broker處理消息過(guò)成功的錯(cuò)誤。

開(kāi)發(fā)者訂閱數(shù)據(jù)時(shí),需要實(shí)現(xiàn)Handler這個(gè)函數(shù),接收Event的實(shí)例,提取數(shù)據(jù)進(jìn)行處理,根據(jù)不同的Broker,可能還需要調(diào)用Ack(),處理出現(xiàn)錯(cuò)誤時(shí),返回error。

go-micro集成RabbitMQ實(shí)戰(zhàn)

大概了解了Broker的定義之后,再來(lái)看下如何使用go-micro收發(fā)RabbitMQ消息。

啟動(dòng)一個(gè)RabbitMQ

如果你已經(jīng)有一個(gè)RabbitMQ服務(wù)器,請(qǐng)?zhí)^(guò)這個(gè)步驟。

這里介紹一個(gè)使用docker快速啟動(dòng)RabbitMQ的方法,當(dāng)然前提是你得安裝了docker。

執(zhí)行如下命令啟動(dòng)一個(gè)rabbitmq的docker容器:

docker run --name rabbitmq1 -p 5672:5672 -p 15672:15672 -d rabbitmq

然后進(jìn)入容器進(jìn)行一些設(shè)置:

docker exec -it rabbitmq1 /bin/bash

啟動(dòng)管理工具、禁用指標(biāo)采集(會(huì)導(dǎo)致某些API500錯(cuò)誤):

rabbitmq-plugins enable rabbitmq_management
 
cd /etc/rabbitmq/conf.d/
echo management_agent.disable_metrics_collector = false > management_agent.disable_metrics_collector.conf

最后重啟容器:

docker restart rabbitmq1

最后瀏覽器中輸入 http://127.0.0.0:15672 即可訪問(wèn),默認(rèn)用戶名和密碼都是 guest 。

編寫(xiě)收發(fā)函數(shù)

為了方便演示,先來(lái)定義發(fā)布消息和接收消息的函數(shù)。其中發(fā)布函數(shù)使用了go-micro提供的Event類型,還有其它類型也可以提供Publish的功能,這里發(fā)送的數(shù)據(jù)格式是Json字符串。接收消息的函數(shù)名稱可以隨意取,但是參數(shù)和返回值必須符合規(guī)范,也就是下邊代碼中的樣子,這個(gè)函數(shù)也可以是綁定到某個(gè)類型的。

// 定義一個(gè)發(fā)布消息的函數(shù):每隔1秒發(fā)布一條消息
func loopPublish(event micro.Event) {
	for {
		time.Sleep(time.Duration(1) * time.Second)

		curUnix := strconv.FormatInt(time.Now().Unix(), 10)
		msg := "{\"Id\":" + curUnix + ",\"Name\":\"張三\"}"
		event.Publish(context.TODO(), msg)
	}
}

// 定義一個(gè)接收消息的函數(shù):將收到的消息打印出來(lái)
func handle(ctx context.Context, msg interface{}) (err error) {
	defer func() {
		if r := recover(); r != nil {
			err = errors.New(fmt.Sprint(r))
			log.Println(err)
		}
	}()

	b, err := json.Marshal(msg)
	if err != nil {
		log.Println(err)
		return
	}

	log.Println(string(b))
	return
}

編寫(xiě)主體代碼

這里先給出代碼,里面提供了一些注釋,后邊還會(huì)有詳細(xì)介紹。

func main() {
	// RabbitMQ的連接參數(shù)
	rabbitmqUrl := "amqp://guest:guest@127.0.0.1:5672/"
	exchangeName := "amq.topic"
	subcribeTopic := "test"
	queueName := "rabbitmqdemo_test"

	// 默認(rèn)是application/protobuf,這里演示用的是Json,所以要改下
	server.DefaultContentType = "application/json"

	// 創(chuàng)建 RabbitMQ Broker
	b := rabbitmq.NewBroker(
		broker.Addrs(rabbitmqUrl),           // RabbitMQ訪問(wèn)地址,含VHost
		rabbitmq.ExchangeName(exchangeName), // 交換機(jī)的名稱
		rabbitmq.DurableExchange(),          // 消息在Exchange中時(shí)會(huì)進(jìn)行持久化處理
		rabbitmq.PrefetchCount(1),           // 同時(shí)消費(fèi)的最大消息數(shù)量
	)

	// 創(chuàng)建Service,內(nèi)部會(huì)初始化一些東西,必須在NewSubscribeOptions前邊
	service := micro.NewService(
		micro.Broker(b),
	)
	service.Init()

	// 初始化訂閱上下文:這里不是必需的,訂閱會(huì)有默認(rèn)值
	subOpts := broker.NewSubscribeOptions(
		rabbitmq.DurableQueue(),   // 隊(duì)列持久化,消費(fèi)者斷開(kāi)連接后,消息仍然保存到隊(duì)列中
		rabbitmq.RequeueOnError(), // 消息處理函數(shù)返回error時(shí),消息再次入隊(duì)列
		rabbitmq.AckOnSuccess(),   // 消息處理函數(shù)沒(méi)有error返回時(shí),go-micro發(fā)送Ack給RabbitMQ
	)

	// 注冊(cè)訂閱
	micro.RegisterSubscriber(
		subcribeTopic,    // 訂閱的Topic
		service.Server(), // 注冊(cè)到的rpcServer
		handle,           // 消息處理函數(shù)
		server.SubscriberContext(subOpts.Context), // 訂閱上下文,也可以使用默認(rèn)的
		server.SubscriberQueue(queueName),         // 隊(duì)列名稱
	)

	// 發(fā)布事件消息
	event := micro.NewEvent(subcribeTopic, service.Client())
	go loopPublish(event)

	log.Println("Service is running ...")
	if err := service.Run(); err != nil {
		log.Println(err)
	}
}

主要邏輯是:

1、先創(chuàng)建一個(gè)RabbitMQ Broker,它實(shí)現(xiàn)了標(biāo)準(zhǔn)的Broker接口。其中主要的參數(shù)是RabbitMQ的訪問(wèn)地址和RabbitMQ交換機(jī),PrefetchCount是訂閱者(或稱為消費(fèi)者)使用的。

2、然后通過(guò) NewService 創(chuàng)建go-micro服務(wù),并將broker設(shè)置進(jìn)去。這里邊會(huì)初始化很多東西,最核心的是創(chuàng)建一個(gè)rpcServer,并將rpcServer和這個(gè)broker綁定起來(lái)。

3、然后是通過(guò) RegisterSubscriber 注冊(cè)訂閱,這個(gè)注冊(cè)有兩個(gè)層面的功能:一是如果RabbitMQ上還不存在這個(gè)隊(duì)列時(shí)創(chuàng)建隊(duì)列,并訂閱指定topic的消息;二是定義go-micro程序從這個(gè)RabbitMQ隊(duì)列接收數(shù)據(jù)的處理方式。

這里詳細(xì)看下訂閱的參數(shù):

func RegisterSubscriber(topic string, s server.Server, h interface{}, opts ...server.SubscriberOption) error
  • topic:go-micro使用的是Topic模式,發(fā)布者發(fā)送消息的時(shí)候要指定一個(gè)topic,訂閱者根據(jù)需要只接收某個(gè)或某幾個(gè)topic的消息;
  • s:消息從RabbitMQ接收后會(huì)進(jìn)入這個(gè)Server進(jìn)行處理,它是NewService的時(shí)候內(nèi)部創(chuàng)建的;
  • h:使用了上一步創(chuàng)建的接收消息的函數(shù) handle,Server中的方法會(huì)調(diào)用這個(gè)函數(shù);
  • opts 是訂閱的一些選項(xiàng),這里需要指定RabbitMQ隊(duì)列的名稱;另外SubscriberContext定義了訂閱的一些行為,這里DurableQueue設(shè)置RabbitMQ訂閱消息的持久化方式,一般我們都希望消息不丟失,這個(gè)設(shè)置的作用是即使程序與RabbitMQ的連接斷開(kāi),消息也會(huì)保存在RabbitMQ隊(duì)列中;AckOnSuccess和RequeueOnError定義了程序處理消息出現(xiàn)錯(cuò)誤時(shí)的行為,如果handle返回error,消息會(huì)重新返回RabbitMQ,然后再投遞給程序。

4、然后這里為了演示,通過(guò)NewEvent創(chuàng)建了一個(gè)Event,通過(guò)它每隔一秒發(fā)送1條消息。

5、最后通過(guò)service.Run()把這個(gè)程序啟動(dòng)起來(lái)。

辛苦寫(xiě)了半天,看一下這個(gè)程序的運(yùn)行效果:

注意一般發(fā)布者和訂閱者是在不同的程序中,這里只是為了方便演示,才把他們放在一個(gè)程序中。所以如果只是發(fā)布消息,就不需要訂閱的代碼,如果只是訂閱,也不需要發(fā)布消息的代碼,大家使用的時(shí)候根據(jù)需要自己裁剪吧。

go-micro集成RabbitMQ的處理流程

這個(gè)部分來(lái)看一下消息在go-micro和RabbitMQ中是怎么流轉(zhuǎn)的,我畫(huà)了一個(gè)示意圖:

這個(gè)圖有點(diǎn)復(fù)雜,這里詳細(xì)講解下。

首先分成三塊:RabbitMQ、消息發(fā)布部分、消息接收部分,這里用不同的顏色進(jìn)行了區(qū)分。

  • RabbitMQ不是本文的重點(diǎn),就把它看成一個(gè)整體就行了。
  • 消息發(fā)布部分:從生產(chǎn)者程序調(diào)用Event.Publish開(kāi)始,然后調(diào)用Client.Publish,到這里為止,都是在go-micro的核心模塊中進(jìn)行處理;然后再調(diào)用Broker.Publish,這里的Broker是RabbitMQ插件的Broker實(shí)例,從這里開(kāi)始進(jìn)入了RabbiitMQ插件部分,然后再依次通過(guò)RabbitMQ Connection的Publish方法、RabbitMQ Channle的Publish方法,最終發(fā)送到RabbitMQ中。
  • 消息接收部分:Service.Run內(nèi)部會(huì)調(diào)用rpcServer.Start,這個(gè)方法內(nèi)部會(huì)調(diào)用Broker.Subscribe,這個(gè)方法是RabbitMQ插件中定義的,它會(huì)讀取RegisterSubscriber時(shí)的一些RabbitMQ隊(duì)列設(shè)置,然后再依次傳遞到RabbitMQ Connection的Consume方法、RabbitMQ Channel的ConsumeQueue方法,最終連接到RabbitMQ,并在RabbitMQ上設(shè)置好要訂閱的隊(duì)列;這些方法還會(huì)返回一個(gè)類型為amqp.Delivery的Go Channel,Broker.Subscribe不斷的從這個(gè)Go Channel中讀取數(shù)據(jù),然后再發(fā)送到調(diào)用Broker.Subscribe時(shí)傳入的一個(gè)消息處理方法中,這里就是rpcServer.HandleEvnet,消息經(jīng)過(guò)一些處理后再進(jìn)入rpcServer內(nèi)部的路由處理模塊,這里就是route.ProcessMessage,這個(gè)方法內(nèi)部會(huì)根據(jù)當(dāng)前消息的topic查找RegisterSubscriber時(shí)注冊(cè)的訂閱,并最終調(diào)用到當(dāng)時(shí)注冊(cè)的用于接收消息的函數(shù)。

這個(gè)處理過(guò)程還可以劃分為業(yè)務(wù)部分、核心模塊部分和插件部分。

  • 首先創(chuàng)建一個(gè)插件的Broker實(shí)現(xiàn),把它注冊(cè)到核心模塊的rpcServer中;
  • 消息的發(fā)送從業(yè)務(wù)部分進(jìn)入核心模塊部分,再進(jìn)入具體實(shí)現(xiàn)Broker的插件部分;
  • 消息的接收則首先進(jìn)入插件部分,然后再流轉(zhuǎn)到核心模塊部分,再流轉(zhuǎn)到業(yè)務(wù)部分。

從上邊的圖中可以看到消息都需要經(jīng)過(guò)這個(gè)RabbitMQ插件進(jìn)行處理,實(shí)際上可以只使用這個(gè)插件,就能實(shí)現(xiàn)消息的發(fā)送和接收。這個(gè)演示代碼我已經(jīng)提交到了Github,有興趣的同學(xué)可以在文末獲取Github倉(cāng)庫(kù)的地址。

從上邊這些劃分中,我們可以理解到設(shè)計(jì)者的整體設(shè)計(jì)思路,把握關(guān)鍵節(jié)點(diǎn),用好用對(duì),出現(xiàn)問(wèn)題時(shí)可以快速定位。

填的幾個(gè)坑

不能接收其它框架發(fā)布的消息

這個(gè)是因?yàn)閞oute.ProcessMessage查找訂閱時(shí)使用了go-micro專用的一個(gè)頭信息:

// get the subscribers by topic
	subs, ok := router.subscribers[msg.Topic()]

這個(gè)msg.Topic返回的是如下實(shí)例中的topic字段:

	rpcMsg := &rpcMessage{
		topic:       msg.Header["Micro-Topic"],
		contentType: ct,
		payload:     &raw.Frame{Data: msg.Body},
		codec:       cf,
		header:      msg.Header,
		body:        msg.Body,
	}

其它框架不會(huì)有這么一個(gè)頭信息,除非專門(mén)適配go-micro。

因?yàn)槭褂肦abbitMQ的場(chǎng)景下,整個(gè)開(kāi)發(fā)都是圍繞RabbitMQ做的,而且go-micro的處理邏輯沒(méi)有考慮RabbitMQ訂閱可以使用通配符的情況,發(fā)布消息的Topic、接收消息的Topic與Micro-Topic的值匹配時(shí)都是按照是否相等的原則處理的,因此可以用RabbitMQ消息自帶的topic來(lái)設(shè)置這個(gè)消息頭。rabbitmq.rbroker.Subscribe 中接收到消息后,就可以進(jìn)行這個(gè)設(shè)置:

// Messages sent from other frameworks to rabbitmq do not have this header.
		// The 'RoutingKey' in the message can be used as this header.
		// Then the message can be transfered to the subscriber which bind this topic.
		msgTopic := header["Micro-Topic"]
		if msgTopic == "" {
			header["Micro-Topic"] = msg.RoutingKey
		}

這樣go-micro開(kāi)發(fā)的消費(fèi)者程序就能接收其它框架發(fā)布的消息了,其它框架無(wú)需適配。

RabbitMQ重啟后訂閱者和發(fā)布者無(wú)限阻塞

go-micro的RabbitMQ插件底層使用另一個(gè)庫(kù):github.com/streadway/amqp

對(duì)于發(fā)布者,RabbitMQ斷開(kāi)連接時(shí)amqp庫(kù)會(huì)通過(guò)Go Channel同步通知go-micro,然后go-micro可以發(fā)起重新連接。問(wèn)題出現(xiàn)在這個(gè)同步通知上,go-micro的RabbitMQ插件設(shè)置了接收連接和通道的關(guān)閉通知,但是只處理了一個(gè)通知就去重新連接了,這就導(dǎo)致有一個(gè)Go Channel一直阻塞,而這個(gè)阻塞會(huì)導(dǎo)致某個(gè)鎖不能釋放,這個(gè)鎖又是Publish時(shí)候需要的,因此導(dǎo)致發(fā)布者無(wú)限阻塞。解決辦法就是外層增加一個(gè)循環(huán),等所有的通知都收到了,再去做重新連接。

對(duì)于訂閱者,RabbitMQ斷開(kāi)連接時(shí),它會(huì)一直阻塞在某個(gè)Go Channel上,直到它返回一個(gè)值,這個(gè)值代表連接已經(jīng)重新建立,訂閱者可以重建消費(fèi)通道。問(wèn)題也是出現(xiàn)在這個(gè)阻塞的Go Channel上,因?yàn)檫@個(gè)Go Channel在每次收到amqp的關(guān)閉通知時(shí)會(huì)重新賦值,而訂閱者等待的Go Channel可能是之前的舊值,永遠(yuǎn)也不會(huì)返回,訂閱者也就無(wú)限阻塞了。解決辦法呢,就是在select時(shí)增加一個(gè)time.After,讓等待的Go Channel有機(jī)會(huì)更新到新值。

代碼就不貼了,有興趣的可以到Github中去看:https://github.com/go-micro/plugins/commit/9f64710807221f3cc649ba4fe05f75b07c66c00c

關(guān)于這兩個(gè)問(wèn)題的修改已經(jīng)合并到官方倉(cāng)庫(kù)中,大家去get最新的代碼就可以了。

這兩個(gè)坑填了,基本上就能滿足我的需要了。當(dāng)然可能還有其它的坑,比如go-micro的RabbitMQ插件好像沒(méi)有發(fā)布者確認(rèn)的功能,這個(gè)要實(shí)現(xiàn),還得好好想想怎么改。

好了,以上就是本文的主要內(nèi)容。

老規(guī)矩,代碼已經(jīng)上傳到Github,歡迎訪問(wèn):https://github.com/bosima/go-demo/tree/main/go-micro-broker-rabbitmq

到此這篇關(guān)于go-micro集成RabbitMQ實(shí)戰(zhàn)和原理詳解的文章就介紹到這了,更多相關(guān)go-micro集成RabbitMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Golang中的內(nèi)存泄漏你真的理解了嗎

    Golang中的內(nèi)存泄漏你真的理解了嗎

    內(nèi)存泄漏是編程中常見(jiàn)的問(wèn)題,會(huì)對(duì)程序的性能和穩(wěn)定性產(chǎn)生嚴(yán)重影響,本文將深入詳解?Golang?中的內(nèi)存泄漏的原因、檢測(cè)方法以及避免方法,希望對(duì)大家有所幫助
    2023-12-12
  • 深入理解?Go?中的字符串

    深入理解?Go?中的字符串

    這篇文章主要介紹了深入理解?Go?中的字符串,在編程語(yǔ)言中,字符串發(fā)揮著重要的角色。字符串背后的數(shù)據(jù)結(jié)構(gòu)一般有兩種類型,一種在編譯時(shí)指定長(zhǎng)度不能修改,一種具有動(dòng)態(tài)的長(zhǎng)度可以修改,下文更多相關(guān)資料需要的小伙伴可以參考一下
    2022-05-05
  • golang中validator包的使用教程

    golang中validator包的使用教程

    Validator 實(shí)際上是一個(gè)驗(yàn)證工具,屬于 golang 的第三方包,這個(gè)包中使用了各種反射技巧來(lái)提供了各種校驗(yàn)和約束數(shù)據(jù)的方式方法,下面就跟隨小編一起來(lái)學(xué)習(xí)一下validator包的使用吧
    2023-09-09
  • Go語(yǔ)言中的Slice學(xué)習(xí)總結(jié)

    Go語(yǔ)言中的Slice學(xué)習(xí)總結(jié)

    這篇文章主要介紹了Go語(yǔ)言中的Slice學(xué)習(xí)總結(jié),本文講解了Slice的定義、Slice的長(zhǎng)度和容量、Slice是引用類型、Slice引用傳遞發(fā)生“意外”等內(nèi)容,需要的朋友可以參考下
    2014-11-11
  • 淺談go中defer的一個(gè)隱藏功能

    淺談go中defer的一個(gè)隱藏功能

    這篇文章主要介紹了淺談go中defer的一個(gè)隱藏功能,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2019-12-12
  • Go語(yǔ)言中的復(fù)合類型詳細(xì)介紹

    Go語(yǔ)言中的復(fù)合類型詳細(xì)介紹

    這篇文章主要介紹了Go語(yǔ)言中的復(fù)合類型詳細(xì)介紹,復(fù)合類型包括:結(jié)構(gòu)體、數(shù)組、切片、Maps,需要的朋友可以參考下
    2014-10-10
  • golang不到30行代碼實(shí)現(xiàn)依賴注入的方法

    golang不到30行代碼實(shí)現(xiàn)依賴注入的方法

    這篇文章主要介紹了golang不到30行代碼實(shí)現(xiàn)依賴注入的方法,小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2018-07-07
  • golang執(zhí)行命令獲取執(zhí)行結(jié)果狀態(tài)(推薦)

    golang執(zhí)行命令獲取執(zhí)行結(jié)果狀態(tài)(推薦)

    這篇文章主要介紹了golang執(zhí)行命令獲取執(zhí)行結(jié)果狀態(tài)的相關(guān)知識(shí),非常不錯(cuò),具有一定的參考借鑒價(jià)值,需要的朋友參考下吧
    2019-11-11
  • Golang 端口復(fù)用測(cè)試的實(shí)現(xiàn)

    Golang 端口復(fù)用測(cè)試的實(shí)現(xiàn)

    這篇文章主要介紹了Golang 端口復(fù)用測(cè)試的實(shí)現(xiàn),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • 使用Singleflight實(shí)現(xiàn)Golang代碼優(yōu)化

    使用Singleflight實(shí)現(xiàn)Golang代碼優(yōu)化

    有許多方法可以優(yōu)化代碼以提高效率,減少運(yùn)行進(jìn)程就是其中之一,本文我們就來(lái)學(xué)習(xí)一下如何通過(guò)使用一個(gè)Go包Singleflight來(lái)減少重復(fù)進(jìn)程,從而優(yōu)化Go代碼吧
    2023-09-09

最新評(píng)論