網站首頁 編程語言 正文
在go-micro中異步消息的收發是通過Broker這個組件來完成的,底層實現有RabbitMQ、Kafka、Redis等等很多種方式,這篇文章主要介紹go-micro使用RabbitMQ收發數據的方法和原理。
Broker的核心功能
Broker的核心功能是Publish和Subscribe,也就是發布和訂閱。它們的定義是:
Publish(topic string, m *Message, opts ...PublishOption) error Subscribe(topic string, h Handler, opts ...SubscribeOption) (Subscriber, error)
發布
發布第一個參數是topic(主題),用于標識某類消息。
發布的數據是通過Message承載的,其包括消息頭和消息體,定義如下:
type Message struct { Header map[string]string Body []byte }
消息頭是map,也就是一組KV(鍵值對)。
消息體是字節數組,在發送和接收時需要開發者進行編碼和解碼的處理。
訂閱
訂閱的第一個參數也是topic(主題),用于過濾出要接收的消息。
訂閱的數據是通過Handler處理的,Handler是一個函數,其定義如下:
type Handler func(Event) error
其中的參數Event是一個接口,需要具體的Broker來實現,其定義如下:
type Event interface { Topic() string Message() *Message Ack() error Error() error }
- Topic() 用于獲取當前消息的topic,也是發布者發送時的topic。
- Message() 用于獲取消息體,也是發布者發送時的Message,其中包括Header和Body。
- Ack() 用于通知Broker消息已經收到了,Broker可以刪除消息了,可用來保證消息至少被消費一次。
- Error() 用于獲取Broker處理消息過成功的錯誤。
開發者訂閱數據時,需要實現Handler這個函數,接收Event的實例,提取數據進行處理,根據不同的Broker,可能還需要調用Ack(),處理出現錯誤時,返回error。
go-micro集成RabbitMQ實戰
大概了解了Broker的定義之后,再來看下如何使用go-micro收發RabbitMQ消息。
啟動一個RabbitMQ
如果你已經有一個RabbitMQ服務器,請跳過這個步驟。
這里介紹一個使用docker快速啟動RabbitMQ的方法,當然前提是你得安裝了docker。
執行如下命令啟動一個rabbitmq的docker容器:
docker run --name rabbitmq1 -p 5672:5672 -p 15672:15672 -d rabbitmq
然后進入容器進行一些設置:
docker exec -it rabbitmq1 /bin/bash
啟動管理工具、禁用指標采集(會導致某些API500錯誤):
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?即可訪問,默認用戶名和密碼都是 guest 。
編寫收發函數
為了方便演示,先來定義發布消息和接收消息的函數。其中發布函數使用了go-micro提供的Event類型,還有其它類型也可以提供Publish的功能,這里發送的數據格式是Json字符串。接收消息的函數名稱可以隨意取,但是參數和返回值必須符合規范,也就是下邊代碼中的樣子,這個函數也可以是綁定到某個類型的。
// 定義一個發布消息的函數:每隔1秒發布一條消息 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) } } // 定義一個接收消息的函數:將收到的消息打印出來 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 }
編寫主體代碼
這里先給出代碼,里面提供了一些注釋,后邊還會有詳細介紹。
func main() { // RabbitMQ的連接參數 rabbitmqUrl := "amqp://guest:guest@127.0.0.1:5672/" exchangeName := "amq.topic" subcribeTopic := "test" queueName := "rabbitmqdemo_test" // 默認是application/protobuf,這里演示用的是Json,所以要改下 server.DefaultContentType = "application/json" // 創建 RabbitMQ Broker b := rabbitmq.NewBroker( broker.Addrs(rabbitmqUrl), // RabbitMQ訪問地址,含VHost rabbitmq.ExchangeName(exchangeName), // 交換機的名稱 rabbitmq.DurableExchange(), // 消息在Exchange中時會進行持久化處理 rabbitmq.PrefetchCount(1), // 同時消費的最大消息數量 ) // 創建Service,內部會初始化一些東西,必須在NewSubscribeOptions前邊 service := micro.NewService( micro.Broker(b), ) service.Init() // 初始化訂閱上下文:這里不是必需的,訂閱會有默認值 subOpts := broker.NewSubscribeOptions( rabbitmq.DurableQueue(), // 隊列持久化,消費者斷開連接后,消息仍然保存到隊列中 rabbitmq.RequeueOnError(), // 消息處理函數返回error時,消息再次入隊列 rabbitmq.AckOnSuccess(), // 消息處理函數沒有error返回時,go-micro發送Ack給RabbitMQ ) // 注冊訂閱 micro.RegisterSubscriber( subcribeTopic, // 訂閱的Topic service.Server(), // 注冊到的rpcServer handle, // 消息處理函數 server.SubscriberContext(subOpts.Context), // 訂閱上下文,也可以使用默認的 server.SubscriberQueue(queueName), // 隊列名稱 ) // 發布事件消息 event := micro.NewEvent(subcribeTopic, service.Client()) go loopPublish(event) log.Println("Service is running ...") if err := service.Run(); err != nil { log.Println(err) } }
主要邏輯是:
1、先創建一個RabbitMQ Broker,它實現了標準的Broker接口。其中主要的參數是RabbitMQ的訪問地址和RabbitMQ交換機,PrefetchCount是訂閱者(或稱為消費者)使用的。
2、然后通過 NewService 創建go-micro服務,并將broker設置進去。這里邊會初始化很多東西,最核心的是創建一個rpcServer,并將rpcServer和這個broker綁定起來。
3、然后是通過 RegisterSubscriber 注冊訂閱,這個注冊有兩個層面的功能:一是如果RabbitMQ上還不存在這個隊列時創建隊列,并訂閱指定topic的消息;二是定義go-micro程序從這個RabbitMQ隊列接收數據的處理方式。
這里詳細看下訂閱的參數:
func RegisterSubscriber(topic string, s server.Server, h interface{}, opts ...server.SubscriberOption) error
- topic:go-micro使用的是Topic模式,發布者發送消息的時候要指定一個topic,訂閱者根據需要只接收某個或某幾個topic的消息;
- s:消息從RabbitMQ接收后會進入這個Server進行處理,它是NewService的時候內部創建的;
- h:使用了上一步創建的接收消息的函數 handle,Server中的方法會調用這個函數;
- opts 是訂閱的一些選項,這里需要指定RabbitMQ隊列的名稱;另外SubscriberContext定義了訂閱的一些行為,這里DurableQueue設置RabbitMQ訂閱消息的持久化方式,一般我們都希望消息不丟失,這個設置的作用是即使程序與RabbitMQ的連接斷開,消息也會保存在RabbitMQ隊列中;AckOnSuccess和RequeueOnError定義了程序處理消息出現錯誤時的行為,如果handle返回error,消息會重新返回RabbitMQ,然后再投遞給程序。
4、然后這里為了演示,通過NewEvent創建了一個Event,通過它每隔一秒發送1條消息。
5、最后通過service.Run()把這個程序啟動起來。
辛苦寫了半天,看一下這個程序的運行效果:
注意一般發布者和訂閱者是在不同的程序中,這里只是為了方便演示,才把他們放在一個程序中。所以如果只是發布消息,就不需要訂閱的代碼,如果只是訂閱,也不需要發布消息的代碼,大家使用的時候根據需要自己裁剪吧。
go-micro集成RabbitMQ的處理流程
這個部分來看一下消息在go-micro和RabbitMQ中是怎么流轉的,我畫了一個示意圖:
這個圖有點復雜,這里詳細講解下。
首先分成三塊:RabbitMQ、消息發布部分、消息接收部分,這里用不同的顏色進行了區分。
- RabbitMQ不是本文的重點,就把它看成一個整體就行了。
- 消息發布部分:從生產者程序調用Event.Publish開始,然后調用Client.Publish,到這里為止,都是在go-micro的核心模塊中進行處理;然后再調用Broker.Publish,這里的Broker是RabbitMQ插件的Broker實例,從這里開始進入了RabbiitMQ插件部分,然后再依次通過RabbitMQ Connection的Publish方法、RabbitMQ Channle的Publish方法,最終發送到RabbitMQ中。
- 消息接收部分:Service.Run內部會調用rpcServer.Start,這個方法內部會調用Broker.Subscribe,這個方法是RabbitMQ插件中定義的,它會讀取RegisterSubscriber時的一些RabbitMQ隊列設置,然后再依次傳遞到RabbitMQ Connection的Consume方法、RabbitMQ Channel的ConsumeQueue方法,最終連接到RabbitMQ,并在RabbitMQ上設置好要訂閱的隊列;這些方法還會返回一個類型為amqp.Delivery的Go Channel,Broker.Subscribe不斷的從這個Go Channel中讀取數據,然后再發送到調用Broker.Subscribe時傳入的一個消息處理方法中,這里就是rpcServer.HandleEvnet,消息經過一些處理后再進入rpcServer內部的路由處理模塊,這里就是route.ProcessMessage,這個方法內部會根據當前消息的topic查找RegisterSubscriber時注冊的訂閱,并最終調用到當時注冊的用于接收消息的函數。
這個處理過程還可以劃分為業務部分、核心模塊部分和插件部分。
- 首先創建一個插件的Broker實現,把它注冊到核心模塊的rpcServer中;
- 消息的發送從業務部分進入核心模塊部分,再進入具體實現Broker的插件部分;
- 消息的接收則首先進入插件部分,然后再流轉到核心模塊部分,再流轉到業務部分。
從上邊的圖中可以看到消息都需要經過這個RabbitMQ插件進行處理,實際上可以只使用這個插件,就能實現消息的發送和接收。這個演示代碼我已經提交到了Github,有興趣的同學可以在文末獲取Github倉庫的地址。
從上邊這些劃分中,我們可以理解到設計者的整體設計思路,把握關鍵節點,用好用對,出現問題時可以快速定位。
填的幾個坑
不能接收其它框架發布的消息
這個是因為route.ProcessMessage查找訂閱時使用了go-micro專用的一個頭信息:
// get the subscribers by topic subs, ok := router.subscribers[msg.Topic()]
這個msg.Topic返回的是如下實例中的topic字段:
rpcMsg := &rpcMessage{ topic: msg.Header["Micro-Topic"], contentType: ct, payload: &raw.Frame{Data: msg.Body}, codec: cf, header: msg.Header, body: msg.Body, }
其它框架不會有這么一個頭信息,除非專門適配go-micro。
因為使用RabbitMQ的場景下,整個開發都是圍繞RabbitMQ做的,而且go-micro的處理邏輯沒有考慮RabbitMQ訂閱可以使用通配符的情況,發布消息的Topic、接收消息的Topic與Micro-Topic的值匹配時都是按照是否相等的原則處理的,因此可以用RabbitMQ消息自帶的topic來設置這個消息頭。rabbitmq.rbroker.Subscribe 中接收到消息后,就可以進行這個設置:
// 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開發的消費者程序就能接收其它框架發布的消息了,其它框架無需適配。
RabbitMQ重啟后訂閱者和發布者無限阻塞
go-micro的RabbitMQ插件底層使用另一個庫:github.com/streadway/amqp
對于發布者,RabbitMQ斷開連接時amqp庫會通過Go Channel同步通知go-micro,然后go-micro可以發起重新連接。問題出現在這個同步通知上,go-micro的RabbitMQ插件設置了接收連接和通道的關閉通知,但是只處理了一個通知就去重新連接了,這就導致有一個Go Channel一直阻塞,而這個阻塞會導致某個鎖不能釋放,這個鎖又是Publish時候需要的,因此導致發布者無限阻塞。解決辦法就是外層增加一個循環,等所有的通知都收到了,再去做重新連接。
對于訂閱者,RabbitMQ斷開連接時,它會一直阻塞在某個Go Channel上,直到它返回一個值,這個值代表連接已經重新建立,訂閱者可以重建消費通道。問題也是出現在這個阻塞的Go Channel上,因為這個Go Channel在每次收到amqp的關閉通知時會重新賦值,而訂閱者等待的Go Channel可能是之前的舊值,永遠也不會返回,訂閱者也就無限阻塞了。解決辦法呢,就是在select時增加一個time.After,讓等待的Go Channel有機會更新到新值。
代碼就不貼了,有興趣的可以到Github中去看:https://github.com/go-micro/plugins/commit/9f64710807221f3cc649ba4fe05f75b07c66c00c
關于這兩個問題的修改已經合并到官方倉庫中,大家去get最新的代碼就可以了。
這兩個坑填了,基本上就能滿足我的需要了。當然可能還有其它的坑,比如go-micro的RabbitMQ插件好像沒有發布者確認的功能,這個要實現,還得好好想想怎么改。
好了,以上就是本文的主要內容。
老規矩,代碼已經上傳到Github,歡迎訪問:https://github.com/bosima/go-demo/tree/main/go-micro-broker-rabbitmq
原文鏈接:https://www.cnblogs.com/bossma/p/16240950.html
相關推薦
- 2022-09-15 C/C++如何實現循環左移,循環右移_C 語言
- 2024-03-07 MyBatis動態語句
- 2022-06-28 Python技法之如何用re模塊實現簡易tokenizer_python
- 2022-10-30 Python接口傳輸url與flask數據詳解_python
- 2022-04-17 使用docker-compose構建鏡像并構建服務時,想為構建的鏡像統一加上指定版本
- 2022-05-25 <C++>詳解類對象作為類成員時調用構造和析構的時機及靜態成員解釋
- 2023-02-23 Go?routine使用方法講解_Golang
- 2022-06-17 Ruby3多線程并行Ractor使用方法詳解_ruby專題
- 最近更新
-
- window11 系統安裝 yarn
- 超詳細win安裝深度學習環境2025年最新版(
- Linux 中運行的top命令 怎么退出?
- MySQL 中decimal 的用法? 存儲小
- get 、set 、toString 方法的使
- @Resource和 @Autowired注解
- Java基礎操作-- 運算符,流程控制 Flo
- 1. Int 和Integer 的區別,Jav
- spring @retryable不生效的一種
- Spring Security之認證信息的處理
- Spring Security之認證過濾器
- Spring Security概述快速入門
- Spring Security之配置體系
- 【SpringBoot】SpringCache
- Spring Security之基于方法配置權
- redisson分布式鎖中waittime的設
- maven:解決release錯誤:Artif
- restTemplate使用總結
- Spring Security之安全異常處理
- MybatisPlus優雅實現加密?
- Spring ioc容器與Bean的生命周期。
- 【探索SpringCloud】服務發現-Nac
- Spring Security之基于HttpR
- Redis 底層數據結構-簡單動態字符串(SD
- arthas操作spring被代理目標對象命令
- Spring中的單例模式應用詳解
- 聊聊消息隊列,發送消息的4種方式
- bootspring第三方資源配置管理
- GIT同步修改后的遠程分支