RabbitMQ延時(shí)消息隊(duì)列在golang中的使用詳解
延時(shí)隊(duì)列常使用在某些業(yè)務(wù)場(chǎng)景,例如訂單支付超時(shí)、接收到外賣(mài)后自動(dòng)確認(rèn)完成訂單、定時(shí)任務(wù)、促銷(xiāo)過(guò)期等,使用延時(shí)隊(duì)列可以簡(jiǎn)化系統(tǒng)的設(shè)計(jì)和開(kāi)發(fā)、提高系統(tǒng)的可靠性和可用性、提高系統(tǒng)的性能。下面介紹使用RabbitMQ的延時(shí)消息隊(duì)列,使用之前先要讓RabbitMQ支持延時(shí)隊(duì)列。
在docker安裝單機(jī)版rabbitMQ
docker-compose.yaml配置文件內(nèi)容如下:
version: '3' services: rabbitmq: image: rabbitmq:3.12-management container_name: rabbitmq hostname: rabbitmq-service restart: always ports: - 5672:5672 - 15672:15672 volumes: - $PWD/data:/var/lib/rabbitmq - $PWD/plugins/enabled_plugins:/etc/rabbitmq/enabled_plugins - $PWD/plugins/rabbitmq_delayed_message_exchange-3.12.0.ez:/plugins/rabbitmq_delayed_message_exchange-3.12.0.ez environment: TZ: Asia/Shanghai RABBITMQ_DEFAULT_USER: guest RABBITMQ_DEFAULT_PASS: guest RABBITMQ_DEFAULT_VHOST: /
rabbitMQ默認(rèn)不支持延時(shí)消息隊(duì)列類(lèi)型,需要另外安裝插件來(lái)實(shí)現(xiàn):
- enabled_plugins 是設(shè)置默認(rèn)開(kāi)啟的插件,內(nèi)容為
[rabbitmq_delayed_message_exchange,rabbitmq_management,rabbitmq_prometheus]
- rabbitmq_delayed_message_exchange-3.12.0.ez 是延時(shí)隊(duì)列插件。
啟動(dòng)rabbitmq:
docker-compose up -d
可以在瀏覽器訪問(wèn)管理后臺(tái) http://localhost:15672 ,用戶名和密碼都是guest
。
點(diǎn)擊菜單【exchange】--> 【Add a new exchange】-->【Type】,在下拉列表中看到x-delayed-message
類(lèi)型的話,說(shuō)明已經(jīng)支持延時(shí)隊(duì)列了。
使用延時(shí)隊(duì)列需要指定具體某一種消息類(lèi)型(direct、topic、fanout、headers),下面以direct類(lèi)型的延時(shí)消息隊(duì)列為例。
生產(chǎn)端示例代碼
package main import ( "context" "fmt" "strconv" "time" amqp "github.com/rabbitmq/amqp091-go" ) var ( url = "amqp://guest:guest@127.0.0.1:5672/" exchangeName = "delayed-message-exchange-demo" ) func main() { conn, err := amqp.Dial(url) checkErr(err) defer conn.Close() ctx := context.Background() queueName := "delayed-message-queue" routingKey := "delayed-key" delayedMessageType := "direct" exchange := NewDelayedMessageExchange(exchangeName, delayedMessageType, routingKey) q, err := NewProducer(queueName, conn, exchange) checkErr(err) defer q.Close() for i := 1; i <= 5; i++ { body := time.Now().Format("2006-01-02 15:04:05.000") + " hello world " + strconv.Itoa(i) err = q.Publish(ctx, time.Second*5, []byte(body)) // 發(fā)送消息 checkErr(err) time.Sleep(time.Second) } } // Exchange 交換機(jī) type Exchange struct { Name string // exchange名稱 Type string // exchange類(lèi)型,支持direct、topic、fanout、headers、x-delayed-message RoutingKey string // 路由key XDelayedMessageType string // 延時(shí)消息類(lèi)型,支持direct、topic、fanout、headers } // NewDelayedMessageExchange 實(shí)例化一個(gè)delayed-message類(lèi)型交換機(jī),參數(shù)delayedMessageType 消息類(lèi)型direct、topic、fanout、headers func NewDelayedMessageExchange(exchangeName string, delayedMessageType string, routingKey string) *Exchange { return &Exchange{ Name: exchangeName, Type: "x-delayed-message", RoutingKey: routingKey, XDelayedMessageType: delayedMessageType, } } // Producer 生產(chǎn)者對(duì)象 type Producer struct { queueName string exchange *Exchange conn *amqp.Connection ch *amqp.Channel } // NewProducer 實(shí)例化一個(gè)生產(chǎn)者 func NewProducer(queueName string, conn *amqp.Connection, exchange *Exchange) (*Producer, error) { // 創(chuàng)建管道 ch, err := conn.Channel() if err != nil { return nil, err } // 聲明交換機(jī)類(lèi)型 err = ch.ExchangeDeclare( exchange.Name, // 交換機(jī)名稱 exchange.Type, // x-delayed-message true, // 是否持久化 false, // 是否自動(dòng)刪除 false, // 是否公開(kāi),false即公開(kāi) false, // 是否等待 amqp.Table{ "x-delayed-type": exchange.XDelayedMessageType, // 延時(shí)消息的類(lèi)型direct、topic、fanout、headers }, ) if err != nil { _ = ch.Close() return nil, err } // 聲明隊(duì)列,如果隊(duì)列不存在則自動(dòng)創(chuàng)建,存在則跳過(guò)創(chuàng)建 q, err := ch.QueueDeclare( queueName, // 消息隊(duì)列名稱 true, // 是否持久化 false, // 是否自動(dòng)刪除 false, // 是否具有排他性(僅創(chuàng)建它的程序才可用) false, // 是否阻塞處理 nil, // 額外的屬性 ) if err != nil { _ = ch.Close() return nil, err } // 綁定隊(duì)列和交換機(jī) err = ch.QueueBind( q.Name, exchange.RoutingKey, exchange.Name, false, nil, ) if err != nil { _ = ch.Close() return nil, err } return &Producer{ queueName: queueName, conn: conn, ch: ch, exchange: exchange, }, nil } // Publish 發(fā)送消息 func (p *Producer) Publish(ctx context.Context, delayTime time.Duration, body []byte) error { err := p.ch.PublishWithContext( ctx, p.exchange.Name, // exchange name p.exchange.RoutingKey, // key false, // mandatory 如果為true,根據(jù)自身exchange類(lèi)型和routingKey規(guī)則無(wú)法找到符合條件的隊(duì)列會(huì)把消息返還給發(fā)送者 false, // immediate 如果為true,當(dāng)exchange發(fā)送消息到隊(duì)列后發(fā)現(xiàn)隊(duì)列上沒(méi)有消費(fèi)者,則會(huì)把消息返還給發(fā)送者 amqp.Publishing{ DeliveryMode: amqp.Persistent, // 如果隊(duì)列的聲明是持久化的,那么消息也設(shè)置為持久化 ContentType: "text/plain", Body: body, Headers: amqp.Table{ "x-delay": int(delayTime / time.Millisecond), // 延遲時(shí)間: 毫秒 }, }, ) if err != nil { return err } fmt.Printf("[send]: %s\n", body) return nil } // Close 關(guān)閉生產(chǎn)者 func (p *Producer) Close() { if p.ch != nil { _ = p.ch.Close() } } func checkErr(err error) { if err != nil { panic(err) } }
消費(fèi)端示例代碼
package main import ( "context" "fmt" "os" "os/signal" "syscall" "time" amqp "github.com/rabbitmq/amqp091-go" ) var ( url = "amqp://guest:guest@127.0.0.1:5672/" exchangeName = "delayed-message-exchange-demo" ) func main() { conn, err := amqp.Dial(url) checkErr(err) defer conn.Close() ctx := context.Background() queueName := "delayed-message-queue" routingKey := "delayed-key" delayedMessageType := "direct" exchange := NewDelayedMessageExchange(exchangeName, delayedMessageType, routingKey) c, err := NewConsumer(ctx, queueName, exchange, conn) checkErr(err) c.Consume() // 消費(fèi)消息 defer c.Close() fmt.Println("exit press CTRL+C") interrupt := make(chan os.Signal, 1) signal.Notify(interrupt, syscall.SIGHUP, syscall.SIGINT, syscall.SIGTERM, syscall.SIGQUIT) <-interrupt fmt.Println("exit consume messages") } // Exchange 交換機(jī) type Exchange struct { Name string // exchange名稱 Type string // exchange類(lèi)型,支持direct、topic、fanout、headers、x-delayed-message RoutingKey string // 路由key XDelayedMessageType string // 延時(shí)消息類(lèi)型,支持direct、topic、fanout、headers } // NewDelayedMessageExchange 實(shí)例化一個(gè)delayed-message類(lèi)型交換機(jī),參數(shù)delayedMessageType 消息類(lèi)型direct、topic、fanout、headers func NewDelayedMessageExchange(exchangeName string, delayedMessageType string, routingKey string) *Exchange { return &Exchange{ Name: exchangeName, Type: "x-delayed-message", RoutingKey: routingKey, XDelayedMessageType: delayedMessageType, } } // Consumer 消費(fèi)者 type Consumer struct { ctx context.Context queueName string conn *amqp.Connection ch *amqp.Channel delivery <-chan amqp.Delivery exchange *Exchange } // NewConsumer 實(shí)例化一個(gè)消費(fèi)者 func NewConsumer(ctx context.Context, queueName string, exchange *Exchange, conn *amqp.Connection) (*Consumer, error) { // 創(chuàng)建管道 ch, err := conn.Channel() if err != nil { return nil, err } // 聲明交換機(jī)類(lèi)型 err = ch.ExchangeDeclare( exchange.Name, // 交換機(jī)名稱 exchange.Type, // 交換機(jī)的類(lèi)型,支持direct、topic、fanout、headers true, // 是否持久化 false, // 是否自動(dòng)刪除 false, // 是否公開(kāi),false即公開(kāi) false, // 是否等待 amqp.Table{ "x-delayed-type": exchange.XDelayedMessageType, // 延時(shí)消息的類(lèi)型direct、topic、fanout、headers }, ) if err != nil { _ = ch.Close() return nil, err } // 聲明隊(duì)列,如果隊(duì)列不存在則自動(dòng)創(chuàng)建,存在則跳過(guò)創(chuàng)建 q, err := ch.QueueDeclare( queueName, // 消息隊(duì)列名稱 true, // 是否持久化 false, // 是否自動(dòng)刪除 false, // 是否具有排他性(僅創(chuàng)建它的程序才可用) false, // 是否阻塞處理 nil, // 額外的屬性 ) if err != nil { _ = ch.Close() return nil, err } // 綁定隊(duì)列和交換機(jī) err = ch.QueueBind( q.Name, exchange.RoutingKey, exchange.Name, false, nil, ) if err != nil { _ = ch.Close() return nil, err } // 為消息隊(duì)列注冊(cè)消費(fèi)者 delivery, err := ch.ConsumeWithContext( ctx, queueName, // queue 名稱 "", // consumer 用來(lái)區(qū)分多個(gè)消費(fèi)者 true, // auto-ack 是否自動(dòng)應(yīng)答 false, // exclusive 是否獨(dú)有 false, // no-local 如果設(shè)置為true,表示不能將同一個(gè)Connection中生產(chǎn)者發(fā)送的消息傳遞給這個(gè)Connection中的消費(fèi)者 false, // no-wait 是否阻塞 nil, // args ) if err != nil { _ = ch.Close() return nil, err } return &Consumer{ queueName: queueName, conn: conn, ch: ch, delivery: delivery, exchange: exchange, }, nil } // Consume 接收消息 func (c *Consumer) Consume() { go func() { fmt.Printf("waiting for messages, type=%s, queue=%s, key=%s\n", c.exchange.Type, c.queueName, c.exchange.RoutingKey) for d := range c.delivery { // 處理消息 fmt.Printf("%s %s [received]: %s\n", time.Now().Format("2006-01-02 15:04:05.000"), c.exchange.RoutingKey, d.Body) // _ = d.Ack(false) // 如果auto-ack為false時(shí),需要手動(dòng)ack } }() } // Close 關(guān)閉 func (c *Consumer) Close() { if c.ch != nil { _ = c.ch.Close() } } func checkErr(err error) { if err != nil { panic(err) } }
總結(jié)
上面介紹了rabbitMQ延時(shí)消息隊(duì)列簡(jiǎn)單使用示例,在實(shí)際使用中,連接rabbitMQ應(yīng)該有網(wǎng)絡(luò)斷開(kāi)重連功能。
rabbitMQ需要依賴插件rabbitmq_delayed_message_exchange,目前該插件的當(dāng)前設(shè)計(jì)并不真正適合包含大量延遲消息(例如數(shù)十萬(wàn)以上)的場(chǎng)景,另外該插件的一個(gè)可變性來(lái)源是依賴于 Erlang 計(jì)時(shí)器,在系統(tǒng)中使用了一定數(shù)量的長(zhǎng)時(shí)間計(jì)時(shí)器之后,它們開(kāi)始爭(zhēng)用調(diào)度程序資源,并且時(shí)間漂移不斷累積。
如果你采用了 Delayed Message 插件這種方式來(lái)實(shí)現(xiàn),對(duì)于消息可靠性要求非常高,在發(fā)送消息之前可以先保存到 DB 打標(biāo)記,消費(fèi)之后將消息標(biāo)記為已消費(fèi),中間可以加入定時(shí)任務(wù)做檢測(cè),這可以進(jìn)一步保證你的消息的可靠性。
這是在github.com/rabbitmq/amqp091-go
基礎(chǔ)上封裝的 rabbitmq 庫(kù),開(kāi)箱即用各種消息類(lèi)型(direct
, topic
, fanout
, headers
, delayed message
, publisher subscriber
)。
到此這篇關(guān)于RabbitMQ延時(shí)消息隊(duì)列在golang中的使用詳解的文章就介紹到這了,更多相關(guān)go RabbitMQ延時(shí)隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Go Excelize API源碼閱讀GetPageLayout及SetPageMargins
這篇文章主要為大家介紹了Go Excelize API源碼閱讀GetPageLayout及SetPageMargins的方法示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-08-08關(guān)于go語(yǔ)言編碼需要放到src 文件夾下的問(wèn)題
這篇文章主要介紹了go語(yǔ)言編碼需要放到src 文件夾下的相關(guān)知識(shí),本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-10-10Go空結(jié)構(gòu)體struct{}的作用是什么
本文主要介紹了Go空結(jié)構(gòu)體struct{}的作用是什么,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2023-02-02Go 語(yǔ)言 IDE 中的 VSCode 配置使用教程
Gogland 是 JetBrains 公司推出的Go語(yǔ)言集成開(kāi)發(fā)環(huán)境。這篇文章主要介紹了Go 語(yǔ)言 IDE 中的 VSCode 配置使用教程,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-05-05Go+Redis實(shí)現(xiàn)常見(jiàn)限流算法的示例代碼
限流是項(xiàng)目中經(jīng)常需要使用到的一種工具,一般用于限制用戶的請(qǐng)求的頻率,也可以避免瞬間流量過(guò)大導(dǎo)致系統(tǒng)崩潰,或者穩(wěn)定消息處理速率。這篇文章主要是使用Go+Redis實(shí)現(xiàn)常見(jiàn)的限流算法,需要的可以參考一下2023-04-04使用Golang編寫(xiě)一個(gè)簡(jiǎn)單的命令行工具
Cobra是一個(gè)強(qiáng)大的開(kāi)源工具,能夠幫助我們快速構(gòu)建出優(yōu)雅且功能豐富的命令行應(yīng)用,本文將利用Cobra編寫(xiě)一個(gè)簡(jiǎn)單的命令行工具,感興趣的可以了解下2023-12-12Go語(yǔ)言實(shí)現(xiàn)互斥鎖、隨機(jī)數(shù)、time、List
這篇文章主要介紹了Go語(yǔ)言實(shí)現(xiàn)互斥鎖、隨機(jī)數(shù)、time、List的相關(guān)資料,需要的朋友可以參考下2018-10-10Go-家庭收支記賬軟件項(xiàng)目實(shí)現(xiàn)
這篇文章主要介紹了Go-家庭收支記賬軟件項(xiàng)目實(shí)現(xiàn),本文章內(nèi)容詳細(xì),具有很好的參考價(jià)值,希望對(duì)大家有所幫助,需要的朋友可以參考下2023-01-01golang?gorm錯(cuò)誤處理事務(wù)以及日志用法示例
這篇文章主要為大家介紹了golang?gorm錯(cuò)誤處理事務(wù)以及日志用法示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步早日升職加薪2022-04-04