在go-micro中異步消息的收發(fā)是通過Broker這個(gè)組件來完成的,底層實(shí)現(xiàn)有RabbitMQ、Kafka、Redis等等很多種方式,這篇文章主要介紹go-micro使用RabbitMQ收發(fā)數(shù)據(jù)的方法和原理。
讓客戶滿意是我們工作的目標(biāo),不斷超越客戶的期望值來自于我們對這個(gè)行業(yè)的熱愛。我們立志把好的技術(shù)通過有效、簡單的方式提供給客戶,將通過不懈努力成為客戶在信息化領(lǐng)域值得信任、有價(jià)值的長期合作伙伴,公司提供的服務(wù)項(xiàng)目有:域名申請、虛擬空間、營銷軟件、網(wǎng)站建設(shè)、松溪網(wǎng)站維護(hù)、網(wǎng)站推廣。
Broker的核心功能是Publish和Subscribe,也就是發(fā)布和訂閱。它們的定義是:
Publish(topic string, m *Message, opts ...PublishOption) error
Subscribe(topic string, h Handler, opts ...SubscribeOption) (Subscriber, error)
發(fā)布第一個(gè)參數(shù)是topic(主題),用于標(biāo)識某類消息。
發(fā)布的數(shù)據(jù)是通過Message承載的,其包括消息頭和消息體,定義如下:
type Message struct {
Header map[string]string
Body []byte
}
消息頭是map,也就是一組KV(鍵值對)。
消息體是字節(jié)數(shù)組,在發(fā)送和接收時(shí)需要開發(fā)者進(jìn)行編碼和解碼的處理。
訂閱的第一個(gè)參數(shù)也是topic(主題),用于過濾出要接收的消息。
訂閱的數(shù)據(jù)是通過Handler處理的,Handler是一個(gè)函數(shù),其定義如下:
type Handler func(Event) error
其中的參數(shù)Event是一個(gè)接口,需要具體的Broker來實(shí)現(xiàn),其定義如下:
type Event interface {
Topic() string
Message() *Message
Ack() error
Error() error
}
開發(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。
大概了解了Broker的定義之后,再來看下如何使用go-micro收發(fā)RabbitMQ消息。
如果你已經(jīng)有一個(gè)RabbitMQ服務(wù)器,請?zhí)^這個(gè)步驟。
這里介紹一個(gè)使用docker快速啟動(dòng)RabbitMQ的方法,當(dāng)然前提是你得安裝了docker。
執(zhí)行如下命令啟動(dòng)一個(gè)rabbitmq的docker容器:
docker run --name rabbitmq1 -p 5672:5672 -p : -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: 即可訪問,默認(rèn)用戶名和密碼都是 guest 。
為了方便演示,先來定義發(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ù):將收到的消息打印出來
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
}
這里先給出代碼,里面提供了一些注釋,后邊還會(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訪問地址,含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)者斷開連接后,消息仍然保存到隊(duì)列中
rabbitmq.RequeueOnError(), // 消息處理函數(shù)返回error時(shí),消息再次入隊(duì)列
rabbitmq.AckOnSuccess(), // 消息處理函數(shù)沒有error返回時(shí),go-micro發(fā)送Ack給RabbitMQ
)
// 注冊訂閱
micro.RegisterSubscriber(
subcribeTopic, // 訂閱的Topic
service.Server(), // 注冊到的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的訪問地址和RabbitMQ交換機(jī),PrefetchCount是訂閱者(或稱為消費(fèi)者)使用的。
2、然后通過 NewService 創(chuàng)建go-micro服務(wù),并將broker設(shè)置進(jìn)去。這里邊會(huì)初始化很多東西,最核心的是創(chuàng)建一個(gè)rpcServer,并將rpcServer和這個(gè)broker綁定起來。
3、然后是通過 RegisterSubscriber 注冊訂閱,這個(gè)注冊有兩個(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
4、然后這里為了演示,通過NewEvent創(chuàng)建了一個(gè)Event,通過它每隔一秒發(fā)送1條消息。
5、最后通過service.Run()把這個(gè)程序啟動(dòng)起來。
辛苦寫了半天,看一下這個(gè)程序的運(yùn)行效果:
注意一般發(fā)布者和訂閱者是在不同的程序中,這里只是為了方便演示,才把他們放在一個(gè)程序中。所以如果只是發(fā)布消息,就不需要訂閱的代碼,如果只是訂閱,也不需要發(fā)布消息的代碼,大家使用的時(shí)候根據(jù)需要自己裁剪吧。
這個(gè)部分來看一下消息在go-micro和RabbitMQ中是怎么流轉(zhuǎn)的,我畫了一個(gè)示意圖:
這個(gè)圖有點(diǎn)復(fù)雜,這里詳細(xì)講解下。
首先分成三塊:RabbitMQ、消息發(fā)布部分、消息接收部分,這里用不同的顏色進(jìn)行了區(qū)分。
這個(gè)處理過程還可以劃分為業(yè)務(wù)部分、核心模塊部分和插件部分。
從上邊的圖中可以看到消息都需要經(jīng)過這個(gè)RabbitMQ插件進(jìn)行處理,實(shí)際上可以只使用這個(gè)插件,就能實(shí)現(xiàn)消息的發(fā)送和接收。這個(gè)演示代碼我已經(jīng)提交到了Github,有興趣的同學(xué)可以在文末獲取Github倉庫的地址。
從上邊這些劃分中,我們可以理解到設(shè)計(jì)者的整體設(shè)計(jì)思路,把握關(guān)鍵節(jié)點(diǎn),用好用對,出現(xiàn)問題時(shí)可以快速定位。
這個(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è)頭信息,除非專門適配go-micro。
因?yàn)槭褂肦abbitMQ的場景下,整個(gè)開發(fā)都是圍繞RabbitMQ做的,而且go-micro的處理邏輯沒有考慮RabbitMQ訂閱可以使用通配符的情況,發(fā)布消息的Topic、接收消息的Topic與Micro-Topic的值匹配時(shí)都是按照是否相等的原則處理的,因此可以用RabbitMQ消息自帶的topic來設(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開發(fā)的消費(fèi)者程序就能接收其它框架發(fā)布的消息了,其它框架無需適配。
go-micro的RabbitMQ插件底層使用另一個(gè)庫:github.com/streadway/amqp
對于發(fā)布者,RabbitMQ斷開連接時(shí)amqp庫會(huì)通過Go Channel同步通知go-micro,然后go-micro可以發(fā)起重新連接。問題出現(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ā)布者無限阻塞。解決辦法就是外層增加一個(gè)循環(huán),等所有的通知都收到了,再去做重新連接。
對于訂閱者,RabbitMQ斷開連接時(shí),它會(huì)一直阻塞在某個(gè)Go Channel上,直到它返回一個(gè)值,這個(gè)值代表連接已經(jīng)重新建立,訂閱者可以重建消費(fèi)通道。問題也是出現(xiàn)在這個(gè)阻塞的Go Channel上,因?yàn)檫@個(gè)Go Channel在每次收到amqp的關(guān)閉通知時(shí)會(huì)重新賦值,而訂閱者等待的Go Channel可能是之前的舊值,永遠(yuǎn)也不會(huì)返回,訂閱者也就無限阻塞了。解決辦法呢,就是在select時(shí)增加一個(gè)time.After,讓等待的Go Channel有機(jī)會(huì)更新到新值。
代碼就不貼了,有興趣的可以到Github中去看:https://github.com/go-micro/plugins/commit/9ff3cc649ba4fe05f75b07c66c00c
關(guān)于這兩個(gè)問題的修改已經(jīng)合并到官方倉庫中,大家去get最新的代碼就可以了。
這兩個(gè)坑填了,基本上就能滿足我的需要了。當(dāng)然可能還有其它的坑,比如go-micro的RabbitMQ插件好像沒有發(fā)布者確認(rèn)的功能,這個(gè)要實(shí)現(xiàn),還得好好想想怎么改。
好了,以上就是本文的主要內(nèi)容。
老規(guī)矩,代碼已經(jīng)上傳到Github,歡迎訪問:https://github.com/bosima/go-demo/tree/main/go-micro-broker-rabbitmq
收獲更多架構(gòu)知識,請關(guān)注微信公眾號 螢火架構(gòu)。原創(chuàng)內(nèi)容,轉(zhuǎn)載請注明出處。
分享題目:go-micro集成RabbitMQ實(shí)戰(zhàn)和原理
當(dāng)前路徑:http://chinadenli.net/article8/dsoisop.html
成都網(wǎng)站建設(shè)公司_創(chuàng)新互聯(lián),為您提供外貿(mào)建站、自適應(yīng)網(wǎng)站、ChatGPT、企業(yè)建站、電子商務(wù)、網(wǎng)站建設(shè)
聲明:本網(wǎng)站發(fā)布的內(nèi)容(圖片、視頻和文字)以用戶投稿、用戶轉(zhuǎn)載內(nèi)容為主,如果涉及侵權(quán)請盡快告知,我們將會(huì)在第一時(shí)間刪除。文章觀點(diǎn)不代表本網(wǎng)站立場,如需處理請聯(lián)系客服。電話:028-86922220;郵箱:631063699@qq.com。內(nèi)容未經(jīng)允許不得轉(zhuǎn)載,或轉(zhuǎn)載時(shí)需注明來源: 創(chuàng)新互聯(lián)