go語(yǔ)言使用RocketMq的示例代碼
前言
在 Go 語(yǔ)言中使用 Apache RocketMQ,官方推薦的方式是通過(guò) apache/rocketmq-client-go 這個(gè)客戶端庫(kù)。它是 Apache 官方維護(hù)的 Go SDK,支持 Producer(生產(chǎn)者)、Push Consumer(推模式消費(fèi)者)、Pull Consumer(拉模式消費(fèi)者)等功能
一、下載并啟動(dòng)
來(lái)到rocketmq下載官網(wǎng)
在Binary下載中選擇合適的版本下載

配置環(huán)境變量

進(jìn)入broker.conf并進(jìn)行修改

修改如下:
brokerClusterName = DefaultCluster brokerName = broker-a brokerId = 0 deleteWhen = 04 fileReservedTime = 48 brokerRole = ASYNC_MASTER flushDiskType = ASYNC_FLUSH
在bin目錄下執(zhí)行如下命令即可啟動(dòng)rocketmq
start mqnamesrv.cmd
.\mqbroker.cmd -n 127.0.0.1:9876 -c …\conf\broker.conf

綜上,rocketmq在windows的配置與啟動(dòng)已完成
二、生產(chǎn)者
import (
"context"
"fmt"
"os"
"github.com/apache/rocketmq-client-go/v2"
"github.com/apache/rocketmq-client-go/v2/primitive"
"github.com/apache/rocketmq-client-go/v2/producer"
)
func main() {
// 創(chuàng)建RocketMQ生產(chǎn)者實(shí)例
// 參數(shù)說(shuō)明:
// - producer.WithNsResolver: 設(shè)置NameServer地址解析器,指定NameServer地址為"127.0.0.1:9876"
// - producer.WithRetry: 設(shè)置消息發(fā)送失敗時(shí)的重試次數(shù)為2次
p, err := rocketmq.NewProducer(
producer.WithNsResolver(primitive.NewPassthroughResolver([]string{"127.0.0.1:9876"})),
producer.WithRetry(2),
)
if err != nil {
fmt.Printf("create producer error: %s\n", err.Error())
os.Exit(1)
}
err = p.Start()
if err != nil {
fmt.Printf("start producer error: %s\n", err.Error())
os.Exit(1)
}
// 延遲執(zhí)行生產(chǎn)者關(guān)閉操作,確保在main函數(shù)結(jié)束前關(guān)閉生產(chǎn)者連接并釋放資源
defer p.Shutdown()
// 發(fā)送同步消息到RocketMQ服務(wù)器
// 構(gòu)建消息對(duì)象,指定主題和消息內(nèi)容
msg := &primitive.Message{
Topic: "TestTopic",
Body: []byte("Hello RocketMQ from Go!"),
}
// 使用生產(chǎn)者同步發(fā)送消息
// 參數(shù)說(shuō)明:
// - context.Background(): 使用空上下文進(jìn)行消息發(fā)送
// - msg: 要發(fā)送的消息對(duì)象
// 返回值說(shuō)明:
// - res: 消息發(fā)送結(jié)果,包含消息ID等信息
// - err: 發(fā)送過(guò)程中可能發(fā)生的錯(cuò)誤
res, err := p.SendSync(context.Background(), msg)//異步發(fā)送消息使用的是 SendAsync 方法
if err != nil {
fmt.Printf("send message error: %s\n", err)
} else {
fmt.Printf("send message success: %s\n", res.String())
}
}
三、消費(fèi)者
推模式示例:
import (
"context"
"fmt"
"os"
"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() {
// 創(chuàng)建RocketMQ推模式消費(fèi)者實(shí)例
// 參數(shù)說(shuō)明:
// - consumer.WithGroupName: 設(shè)置消費(fèi)者組名稱為"TestGroup"
// - consumer.WithNsResolver: 設(shè)置NameServer地址解析器,指定NameServer地址為"127.0.0.1:9876"
// 返回值說(shuō)明:
// - c: 創(chuàng)建的PushConsumer實(shí)例
// - err: 創(chuàng)建過(guò)程中可能發(fā)生的錯(cuò)誤
c, err := rocketmq.NewPushConsumer(
consumer.WithGroupName("TestGroup"),
consumer.WithNsResolver(primitive.NewPassthroughResolver([]string{"127.0.0.1:9876"})),
)
if err != nil {
fmt.Printf("create consumer error: %s\n", err.Error())
os.Exit(1)
}
// 訂閱指定主題的消息
// 參數(shù)說(shuō)明:
// - "TestTopic": 要訂閱的主題名稱
// - consumer.MessageSelector{}: 消息選擇器,此處使用默認(rèn)選擇器
// - func(ctx context.Context, msgs ...*primitive.MessageExt) (consumer.ConsumeResult, error): 消息處理回調(diào)函數(shù)
// - ctx: 上下文參數(shù)
// - msgs: 接收到的消息列表
// - 返回值: 消費(fèi)結(jié)果和可能的錯(cuò)誤
// 返回值說(shuō)明:
// - err: 訂閱過(guò)程中可能發(fā)生的錯(cuò)誤
err = c.Subscribe("TestTopic", consumer.MessageSelector{}, func(ctx context.Context,
msgs ...*primitive.MessageExt) (consumer.ConsumeResult, error) {
for _, msg := range msgs {
fmt.Printf("Received message: %s\n", string(msg.Body))
}
return consumer.ConsumeSuccess, nil
})
if err != nil {
fmt.Printf("subscribe error: %s\n", err.Error())
os.Exit(1)
}
err = c.Start()
if err != nil {
fmt.Printf("start consumer error: %s\n", err.Error())
os.Exit(1)
}
defer c.Shutdown()
// 阻塞主 goroutine,防止程序退出
select {}
}
消費(fèi)者有推模式和拉模式兩種模式,如下
1、推模式
優(yōu)點(diǎn):
實(shí)時(shí)性高 消息到達(dá) Broker 后幾乎立即被投遞給消費(fèi)者(延遲低),適合對(duì)實(shí)時(shí)性要求高的場(chǎng)景。
RocketMQ Push Consumer 會(huì)自動(dòng)進(jìn)行 隊(duì)列分配(Rebalance),多個(gè)消費(fèi)者實(shí)例能自動(dòng)分?jǐn)?Topic 的消息隊(duì)列
消費(fèi)成功后自動(dòng)提交 offset;失敗可重試(支持重試 Topic)
缺點(diǎn):
消費(fèi)者被動(dòng)接收,難以控制消費(fèi)速率
消息由 Broker 主動(dòng)“推送”(通過(guò)長(zhǎng)輪詢模擬),消費(fèi)者無(wú)法決定何時(shí)拉、拉多少。
如果消息突發(fā)洪峰,消費(fèi)者可能被瞬間壓垮(即使有流控,也可能來(lái)不及響應(yīng))
常用場(chǎng)景:
實(shí)時(shí)消息處理(如訂單支付通知、即時(shí)通訊)
希望快速消費(fèi)、低延遲
不想手動(dòng)管理拉取邏輯和 offset
2、拉模式
優(yōu)點(diǎn):
消費(fèi)者自己控制何時(shí)、從哪個(gè)隊(duì)列、拉多少條消息,完全手動(dòng)管理消費(fèi)進(jìn)度(offset)。
可精確控制拉取頻率、批量大小、消費(fèi)速度,適合復(fù)雜調(diào)度邏輯。
缺點(diǎn):
開(kāi)發(fā)復(fù)雜度高,易出錯(cuò),必須手動(dòng)管理 offset,無(wú)自動(dòng)負(fù)載均衡
實(shí)時(shí)性差,延遲不可控。若為了低延遲頻繁拉取,又會(huì)增加 Broker 壓力,產(chǎn)生大量空輪詢(浪費(fèi)網(wǎng)絡(luò)和 CPU)
常用場(chǎng)景:
消費(fèi)速度需要嚴(yán)格控制(如限流、批處理)
需要自定義消費(fèi)策略(如按優(yōu)先級(jí)消費(fèi))
消費(fèi)邏輯復(fù)雜,需精細(xì)控制 offset
到此這篇關(guān)于go語(yǔ)言使用RocketMq的示例代碼的文章就介紹到這了,更多相關(guān)go語(yǔ)言使用RocketMq內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
使用client go實(shí)現(xiàn)自定義控制器的方法
本文我們來(lái)使用client-go實(shí)現(xiàn)一個(gè)自定義控制器,通過(guò)判斷service的Annotations屬性是否包含ingress/http,如果包含則創(chuàng)建ingress,如果不包含則不創(chuàng)建,對(duì)client go自定義控制器相關(guān)知識(shí)感興趣的朋友一起看看吧2022-05-05
go語(yǔ)言中time包的各種函數(shù)總結(jié)
時(shí)間和日期是我們編程中經(jīng)常會(huì)用到的,下面這篇文章主要給大家介紹了關(guān)于go語(yǔ)言中time包的各種函數(shù)總結(jié)的相關(guān)資料,文中通過(guò)實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下2023-04-04
Golang實(shí)現(xiàn)短網(wǎng)址/短鏈服務(wù)的開(kāi)發(fā)筆記分享
這篇文章主要為大家詳細(xì)介紹了如何使用Golang實(shí)現(xiàn)短網(wǎng)址/短鏈服務(wù),文中的示例代碼講解詳細(xì),具有一定的學(xué)習(xí)價(jià)值,感興趣的小伙伴可以了解一下2023-05-05
golang協(xié)程池模擬實(shí)現(xiàn)群發(fā)郵件功能
這篇文章主要介紹了golang協(xié)程池模擬實(shí)現(xiàn)群發(fā)郵件功能,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-05-05
Go開(kāi)發(fā)環(huán)境搭建詳細(xì)介紹
由于目前網(wǎng)上Go的開(kāi)發(fā)環(huán)境搭建文章很多,有些比較老舊,都是基于 GOPATH的,給新入門的同學(xué)造成困擾。以下為2023 版 Go 開(kāi)發(fā)環(huán)境搭建,可參照此教程搭建Go開(kāi)發(fā)環(huán)境,有需要的朋友可以參考閱讀2023-04-04

