最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

使用Go快速驗(yàn)證RocketMQ是否正常的完整過(guò)程

 更新時(shí)間:2026年06月28日 13:45:05   作者:XMYX-0  
最近需要驗(yàn)證一套 RocketMQ 環(huán)境是否可用,最開始的訴求很簡(jiǎn)單:我有一個(gè) RocketMQ NameServer 地址,端口是默認(rèn)端口 9876,想確認(rèn)它到底是不是正常,所以本文給大家介紹了如何使用Go快速驗(yàn)證RocketMQ是否正常,需要的朋友可以參考下

前言

最近需要驗(yàn)證一套 RocketMQ 環(huán)境是否可用。最開始的訴求很簡(jiǎn)單:我有一個(gè) RocketMQ NameServer 地址,端口是默認(rèn)端口 9876,想確認(rèn)它到底是不是正常。

很多時(shí)候我們會(huì)先用 telnet、nc 或者簡(jiǎn)單的 TCP 腳本測(cè)一下端口。這個(gè)動(dòng)作當(dāng)然有價(jià)值,但它只能說(shuō)明:

當(dāng)前機(jī)器到 RocketMQ NameServer 的網(wǎng)絡(luò)是通的。

它不能證明 broker 正常、topic 可用,也不能證明消息真的可以成功生產(chǎn)和消費(fèi)。

所以這次我做了一個(gè)小工具,分三層驗(yàn)證 RocketMQ:

  1. TCP 連通性測(cè)試
  2. 發(fā)送一條消息
  3. 發(fā)送并消費(fèi)同一條消息,完成生產(chǎn)消費(fèi)閉環(huán)

下面記錄一下完整過(guò)程。

環(huán)境信息

本文中的內(nèi)網(wǎng)地址做了脫敏處理,實(shí)際使用時(shí)替換成自己的 RocketMQ 地址即可。

NameServer: 10.x.x.x:9876
Go: go1.25.0 linux/amd64
RocketMQ Go Client: github.com/apache/rocketmq-client-go/v2

測(cè)試機(jī)器上 Go 環(huán)境正常:

go version

輸出類似:

go version go1.25.0 linux/amd64

為什么沒(méi)有繼續(xù)用 Python

一開始我嘗試過(guò) Python 版本,先做 TCP 測(cè)試,再用 rocketmq-client-python 做生產(chǎn)消費(fèi)測(cè)試。

TCP 測(cè)試是正常的:

[OK] TCP connected to 10.x.x.x:9876

但是 Python 客戶端遇到了兩個(gè)問(wèn)題。

在 Linux 上:
rocketmq-client-python 依賴 native 動(dòng)態(tài)庫(kù),如果機(jī)器上缺少 librocketmq.so,會(huì)報(bào):

rocketmq dynamic library not found

Windows 上則更直接:

rocketmq-python does not support Windows

所以最后我換成了 Go。Go 版本更適合這種輕量健康檢查腳本,不需要額外處理 Python native 動(dòng)態(tài)庫(kù)的問(wèn)題。

項(xiàng)目結(jié)構(gòu)

新建一個(gè)目錄:

mkdir rocketmq-go-healthcheck
cd rocketmq-go-healthcheck

目錄結(jié)構(gòu)如下:

rocketmq-go-healthcheck/
├── go.mod
├── go.sum
└── main.go

go.mod

module rocketmq-go-healthcheck

go 1.20

require github.com/apache/rocketmq-client-go/v2 v2.1.2

require (
	github.com/emirpasic/gods v1.12.0 // indirect
	github.com/golang/mock v1.3.1 // indirect
	github.com/google/uuid v1.3.0 // indirect
	github.com/json-iterator/go v1.1.12 // indirect
	github.com/konsorten/go-windows-terminal-sequences v1.0.1 // indirect
	github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 // indirect
	github.com/modern-go/reflect2 v1.0.2 // indirect
	github.com/patrickmn/go-cache v2.1.0+incompatible // indirect
	github.com/pkg/errors v0.8.1 // indirect
	github.com/sirupsen/logrus v1.4.0 // indirect
	github.com/tidwall/gjson v1.13.0 // indirect
	github.com/tidwall/match v1.1.1 // indirect
	github.com/tidwall/pretty v1.2.0 // indirect
	go.uber.org/atomic v1.5.1 // indirect
	golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9 // indirect
	golang.org/x/lint v0.0.0-20190930215403-16217165b5de // indirect
	golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f // indirect
	golang.org/x/tools v0.0.0-20201224043029-2b0845dc783e // indirect
	gopkg.in/natefinch/lumberjack.v2 v2.0.0 // indirect
	stathat.com/c/consistent v1.0.0 // indirect
)

然后執(zhí)行:

go mod tidy # 只需要go.mod 和 main.go 然后執(zhí)行這個(gè)會(huì)自動(dòng)生成go.sum的

核心代碼

下面是完整的 main.go

package main

import (
	"context"
	"flag"
	"fmt"
	"net"
	"os"
	"strings"
	"sync"
	"time"

	rocketmq "github.com/apache/rocketmq-client-go/v2"
	"github.com/apache/rocketmq-client-go/v2/consumer"
	"github.com/apache/rocketmq-client-go/v2/primitive"
	"github.com/apache/rocketmq-client-go/v2/producer"
)

func main() {
	var (
		// NameServer 地址,默認(rèn)端口是 9876。
		namesrv = flag.String("namesrv", "10.x.x.x:9876", "RocketMQ NameServer address")

		// 用于測(cè)試的 topic,需要替換成實(shí)際環(huán)境中存在的 topic。
		topic = flag.String("topic", "TopicTest", "topic used for send/roundtrip")

		// 測(cè)試消息使用的 tag,消費(fèi)者也會(huì)按這個(gè) tag 訂閱。
		tag = flag.String("tag", "healthcheck", "message tag")

		// 支持三種模式:tcp 只測(cè)端口,send 只發(fā)消息,roundtrip 發(fā)送并消費(fèi)。
		mode = flag.String("mode", "roundtrip", "tcp, send, or roundtrip")

		producerGroup = flag.String("producer-group", "PID_HEALTHCHECK_GO", "producer group")

		// consumer group 默認(rèn)帶時(shí)間戳,避免多次測(cè)試時(shí)和已有消費(fèi)者互相影響。
		consumerGroup = flag.String("consumer-group", fmt.Sprintf("CID_HEALTHCHECK_GO_%d", time.Now().UnixNano()), "consumer group")

		tcpTimeout  = flag.Duration("tcp-timeout", 3*time.Second, "tcp connection timeout")
		waitTimeout = flag.Duration("wait", 20*time.Second, "roundtrip consume wait timeout")

		// consumer 啟動(dòng)后稍等一下再發(fā)送消息,降低訂閱尚未完成導(dǎo)致誤判的概率。
		warmup = flag.Duration("warmup", 2*time.Second, "consumer warmup time before send")
	)
	flag.Parse()

	fmt.Printf("[INFO] NameServer: %s\n", *namesrv)
	fmt.Printf("[INFO] Mode: %s\n", *mode)

	// 第一層檢查:先測(cè) TCP。
	// 如果這里不通,優(yōu)先排查網(wǎng)絡(luò)、防火墻、安全組、端口等問(wèn)題。
	if err := checkTCP(*namesrv, *tcpTimeout); err != nil {
		fmt.Printf("[FAIL] TCP connection failed: %v\n", err)
		os.Exit(2)
	}
	fmt.Println("[OK] TCP connected")

	switch *mode {
	case "tcp":
		return
	case "send":
		// 第二層檢查:只發(fā)送一條消息,用于驗(yàn)證 producer 到 broker 的鏈路。
		if err := sendMessage(*namesrv, *topic, *tag, *producerGroup, marker()); err != nil {
			fmt.Printf("[FAIL] Send failed: %v\n", err)
			os.Exit(3)
		}
	case "roundtrip":
		// 第三層檢查:發(fā)送并消費(fèi)同一條消息,更接近真實(shí)可用性驗(yàn)證。
		if err := roundtrip(*namesrv, *topic, *tag, *producerGroup, *consumerGroup, *waitTimeout, *warmup); err != nil {
			fmt.Printf("[FAIL] Roundtrip failed: %v\n", err)
			os.Exit(4)
		}
	default:
		fmt.Printf("[FAIL] Unknown mode %q, use tcp, send, or roundtrip\n", *mode)
		os.Exit(1)
	}
}

func checkTCP(address string, timeout time.Duration) error {
	if !strings.Contains(address, ":") {
		address += ":9876"
	}

	start := time.Now()
	conn, err := net.DialTimeout("tcp", address, timeout)
	if err != nil {
		return err
	}
	_ = conn.Close()

	fmt.Printf("[OK] TCP connected to %s in %d ms\n", address, time.Since(start).Milliseconds())
	return nil
}

func sendMessage(namesrv, topic, tag, group, key string) error {
	p, err := rocketmq.NewProducer(
		producer.WithNameServer([]string{namesrv}),
		producer.WithGroupName(group),
	)
	if err != nil {
		return err
	}
	if err := p.Start(); err != nil {
		return err
	}
	defer p.Shutdown()

	// 每條測(cè)試消息都帶唯一 key,方便 roundtrip 模式精確識(shí)別本次發(fā)送的消息。
	body := fmt.Sprintf("rocketmq-go-healthcheck %s", key)
	msg := primitive.NewMessage(topic, []byte(body))
	msg.WithTag(tag)
	msg.WithKeys([]string{key})

	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()

	result, err := p.SendSync(ctx, msg)
	if err != nil {
		return err
	}

	fmt.Printf("[OK] Sent message: status=%v msg_id=%s queue=%s offset=%d\n",
		result.Status, result.MsgID, result.MessageQueue.String(), result.QueueOffset)
	return nil
}

func roundtrip(namesrv, topic, tag, producerGroup, consumerGroup string, waitTimeout, warmup time.Duration) error {
	// 本次健康檢查的唯一標(biāo)識(shí),避免誤消費(fèi) topic 中已有的歷史消息。
	key := marker()
	found := make(chan *primitive.MessageExt, 1)
	var once sync.Once

	c, err := rocketmq.NewPushConsumer(
		consumer.WithNameServer([]string{namesrv}),
		consumer.WithGroupName(consumerGroup),
		consumer.WithConsumerModel(consumer.Clustering),
	)
	if err != nil {
		return err
	}

	// 只訂閱指定 tag 的消息,減少無(wú)關(guān)消息干擾。
	err = c.Subscribe(topic, consumer.MessageSelector{
		Type:       consumer.TAG,
		Expression: tag,
	}, func(ctx context.Context, msgs ...*primitive.MessageExt) (consumer.ConsumeResult, error) {
		for _, msg := range msgs {
			body := string(msg.Body)
			if body == "" {
				continue
			}

			// 只認(rèn)本次發(fā)送的測(cè)試消息。
			if containsKey(msg, key) || strings.Contains(body, key) {
				once.Do(func() {
					found <- msg
				})
			}
		}
		return consumer.ConsumeSuccess, nil
	})
	if err != nil {
		return err
	}

	if err := c.Start(); err != nil {
		return err
	}
	defer c.Shutdown()

	fmt.Printf("[INFO] Consumer started: group=%s topic=%s tag=%s\n", consumerGroup, topic, tag)
	time.Sleep(warmup)

	// consumer 啟動(dòng)并完成短暫預(yù)熱后,再發(fā)送測(cè)試消息。
	if err := sendMessage(namesrv, topic, tag, producerGroup, key); err != nil {
		return err
	}

	select {
	case msg := <-found:
		fmt.Printf("[OK] Consumed healthcheck message: msg_id=%s queue=%s offset=%d body=%s\n",
			msg.MsgId, msg.Queue.String(), msg.QueueOffset, string(msg.Body))
		return nil
	case <-time.After(waitTimeout):
		return fmt.Errorf("message was sent but not consumed within %s", waitTimeout)
	}
}

func marker() string {
	return fmt.Sprintf("healthcheck-%d", time.Now().UnixNano())
}

func containsKey(msg *primitive.MessageExt, key string) bool {
	keys := msg.GetKeys()
	return keys == key || strings.Contains(keys, key)
}

支持的測(cè)試模式

這個(gè)工具支持三種模式。

只測(cè)試 TCP 連通性

go run . -mode tcp -namesrv 10.x.x.x:9876

如果成功,會(huì)看到類似輸出:

[INFO] NameServer: 10.x.x.x:9876
[INFO] Mode: tcp
[OK] TCP connected to 10.x.x.x:9876 in 40 ms
[OK] TCP connected

這個(gè)結(jié)果只能說(shuō)明網(wǎng)絡(luò)和端口沒(méi)問(wèn)題,不能說(shuō)明 RocketMQ 生產(chǎn)消費(fèi)一定正常。

只發(fā)送一條消息

go run . -mode send -namesrv 10.x.x.x:9876 -topic TopicTest

成功輸出類似:

[OK] Sent message: status=SendOK msg_id=xxx queue=xxx offset=123

如果這一步失敗,常見(jiàn)原因包括:

  1. topic 不存在
  2. broker 沒(méi)有正常注冊(cè)到 NameServer
  3. 客戶端到 broker 網(wǎng)絡(luò)不通
  4. ACL 權(quán)限不足

發(fā)送并消費(fèi)閉環(huán)測(cè)試

go run . -namesrv 10.x.x.x:9876 -topic TopicTest

默認(rèn)模式就是 roundtrip,等價(jià)于:

go run . -mode roundtrip -namesrv 10.x.x.x:9876 -topic TopicTest

成功輸出類似:

[INFO] Consumer started: group=CID_HEALTHCHECK_GO_xxx topic=TopicTest tag=healthcheck
[OK] Sent message: status=SendOK msg_id=xxx queue=xxx offset=123
[OK] Consumed healthcheck message: msg_id=xxx queue=xxx offset=123 body=rocketmq-go-healthcheck healthcheck-xxx

看到這類輸出,基本就可以說(shuō)明:

  1. NameServer 可連接
  2. broker 路由可獲取
  3. topic 可用
  4. producer 可以發(fā)送消息
  5. consumer 可以訂閱并消費(fèi)消息

這比單純測(cè)端口更有說(shuō)服力。

幾個(gè)實(shí)現(xiàn)細(xì)節(jié)

為什么先做 TCP 檢查

代碼一開始先調(diào)用 checkTCP。

這樣做的好處是可以快速區(qū)分問(wèn)題類型:

TCP 不通:優(yōu)先查網(wǎng)絡(luò)、防火墻、安全組、端口
TCP 通但發(fā)送失?。簝?yōu)先查 RocketMQ 路由、topic、broker、ACL
發(fā)送成功但消費(fèi)失?。簝?yōu)先查 consumer group、tag、訂閱、broker 到客戶端網(wǎng)絡(luò)

排障時(shí)分層非常重要,不然很容易一上來(lái)就陷入客戶端報(bào)錯(cuò)細(xì)節(jié)里。

為什么 roundtrip 里先啟動(dòng) consumer

代碼里是先啟動(dòng) consumer,再發(fā)送消息:

if err := c.Start(); err != nil {
	return err
}

time.Sleep(warmup)

if err := sendMessage(namesrv, topic, tag, producerGroup, key); err != nil {
	return err
}

這是為了避免消費(fèi)者還沒(méi)完成訂閱初始化,消息就已經(jīng)發(fā)出去了。雖然 RocketMQ 本身是消息隊(duì)列,但對(duì)于這種一次性的健康檢查腳本來(lái)說(shuō),先啟動(dòng) consumer,等一小會(huì)兒,再發(fā)送消息,結(jié)果更穩(wěn)定。

為什么每次消息都帶唯一 marker

每次運(yùn)行會(huì)生成一個(gè)類似這樣的 key:

healthcheck-171xxxxxxxxxxxxx

消費(fèi)回調(diào)里只接受包含這個(gè) marker 的消息:

if containsKey(msg, key) || strings.Contains(body, key) {
	found <- msg
}

這樣可以避免誤把 topic 里其他歷史消息當(dāng)成本次健康檢查結(jié)果。

常見(jiàn)問(wèn)題

TCP 通,但 send 失敗

這種情況說(shuō)明 NameServer 端口能訪問(wèn),但 RocketMQ 鏈路還不完整。

可以重點(diǎn)檢查:

topic 是否存在
broker 是否啟動(dòng)
broker 是否注冊(cè)到了 NameServer
客戶端機(jī)器是否能訪問(wèn) broker 地址
是否開啟了 ACL

尤其要注意:客戶端連接的是 NameServer,但真正發(fā)送消息時(shí),還會(huì)根據(jù)路由連接 broker。所以 NameServer 通,不代表 broker 一定通。

send 成功,但 roundtrip 消費(fèi)不到

可以檢查:

topic 是否正確
tag 是否匹配
consumer group 是否被其他程序占用
等待時(shí)間是否太短
broker 到客戶端網(wǎng)絡(luò)是否正常

可以適當(dāng)把等待時(shí)間調(diào)大:

go run . -topic TopicTest -wait 60s

topic 不是 TopicTest 怎么辦

直接用 -topic 參數(shù)指定真實(shí) topic:

go run . -namesrv 10.x.x.x:9876 -topic YourRealTopic

小結(jié)

這次驗(yàn)證 RocketMQ 的過(guò)程可以總結(jié)成一句話:

TCP 通只能說(shuō)明端口可達(dá),生產(chǎn)消費(fèi)閉環(huán)成功,才能更接近真實(shí)可用。

最終我用 Go 寫了一個(gè)小工具,支持:

  1. tcp:只測(cè) NameServer 端口
  2. send:發(fā)送一條測(cè)試消息
  3. roundtrip:發(fā)送并消費(fèi)同一條測(cè)試消息

實(shí)際排障時(shí)建議按這個(gè)順序來(lái):

go run . -mode tcp -namesrv 10.x.x.x:9876
go run . -mode send -namesrv 10.x.x.x:9876 -topic TopicTest
go run . -mode roundtrip -namesrv 10.x.x.x:9876 -topic TopicTest

這樣能快速判斷問(wèn)題到底在網(wǎng)絡(luò)層、RocketMQ 路由層,還是生產(chǎn)消費(fèi)鏈路上。

以上就是使用Go快速驗(yàn)證RocketMQ是否正常的完整過(guò)程的詳細(xì)內(nèi)容,更多關(guān)于Go驗(yàn)證RocketMQ是否正常的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 詳解Golang time包中的time.Duration類型

    詳解Golang time包中的time.Duration類型

    在日常開發(fā)過(guò)程中,會(huì)頻繁遇到對(duì)時(shí)間進(jìn)行操作的場(chǎng)景,使用 Golang 中的 time 包可以很方便地實(shí)現(xiàn)對(duì)時(shí)間的相關(guān)操作,本文講解一下 time 包中的 time.Duration 類型,需要的朋友可以參考下
    2023-07-07
  • go語(yǔ)言使用Casbin實(shí)現(xiàn)角色的權(quán)限控制

    go語(yǔ)言使用Casbin實(shí)現(xiàn)角色的權(quán)限控制

    Casbin是用于Golang項(xiàng)目的功能強(qiáng)大且高效的開源訪問(wèn)控制庫(kù)。本文主要介紹了go語(yǔ)言使用Casbin實(shí)現(xiàn)角色的權(quán)限控制,感興趣的可以了解下
    2021-06-06
  • GoLang內(nèi)存模型詳細(xì)講解

    GoLang內(nèi)存模型詳細(xì)講解

    go官方介紹go內(nèi)存模型的時(shí)候說(shuō):探究在什么條件下,goroutine 在讀取一個(gè)變量的值的時(shí),能夠看到其它 goroutine 對(duì)這個(gè)變量進(jìn)行的寫的結(jié)果,Go內(nèi)存模型規(guī)定了一些條件,在這些條件下,在一個(gè)goroutine中讀取變量返回的值能夠確保是另一個(gè)goroutine中對(duì)該變量寫入的值
    2022-12-12
  • golang的基礎(chǔ)語(yǔ)法和常用開發(fā)工具詳解

    golang的基礎(chǔ)語(yǔ)法和常用開發(fā)工具詳解

    這篇文章主要介紹了golang的基礎(chǔ)語(yǔ)法和常用開發(fā)工具,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-12-12
  • Golang中的Map是如何實(shí)現(xiàn)遍歷的

    Golang中的Map是如何實(shí)現(xiàn)遍歷的

    Golang中Map遍歷方法解析,介紹for...range循環(huán)遍歷無(wú)序鍵值對(duì)集合map的常見(jiàn)方式,強(qiáng)調(diào)協(xié)程環(huán)境下的同步保護(hù)措施
    2026-06-06
  • Golang學(xué)習(xí)筆記(二):類型、變量、常量

    Golang學(xué)習(xí)筆記(二):類型、變量、常量

    這篇文章主要介紹了Golang學(xué)習(xí)筆記(二):類型、變量、常量,本文講解了基本類型、保留字、變量、常量、枚舉、運(yùn)算符、指針、分組聲明等內(nèi)容,需要的朋友可以參考下
    2015-05-05
  • Go語(yǔ)言單例模式詳解

    Go語(yǔ)言單例模式詳解

    本文主要介紹了Go語(yǔ)言單例模式詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2023-03-03
  • 淺談Go語(yǔ)言多態(tài)的實(shí)現(xiàn)與interface使用

    淺談Go語(yǔ)言多態(tài)的實(shí)現(xiàn)與interface使用

    如果大家系統(tǒng)的學(xué)過(guò)C++、Java等語(yǔ)言以及面向?qū)ο蟮脑?,相信?yīng)該對(duì)多態(tài)不會(huì)陌生。多態(tài)是面向?qū)ο蠓懂牣?dāng)中經(jīng)常使用并且非常好用的一個(gè)功能,它主要是用在強(qiáng)類型語(yǔ)言當(dāng)中,像是Python這樣的弱類型語(yǔ)言,變量的類型可以隨意變化,也沒(méi)有任何限制,其實(shí)區(qū)別不是很大
    2021-06-06
  • Go語(yǔ)言nil標(biāo)識(shí)符(空值/零值)

    Go語(yǔ)言nil標(biāo)識(shí)符(空值/零值)

    本文主要介紹了Go語(yǔ)言nil標(biāo)識(shí)符(空值/零值),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2023-03-03
  • 如何在 ubuntu linux 上配置 go 語(yǔ)言的 qt 開發(fā)環(huán)境

    如何在 ubuntu linux 上配置 go 語(yǔ)言的 qt 開發(fā)環(huán)境

    這篇文章主要介紹了如何在 ubuntu linux 上配置 go 語(yǔ)言的 qt 開發(fā)環(huán)境,本文分步驟通過(guò)實(shí)例代碼相結(jié)合給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-04-04

最新評(píng)論

永修县| 临高县| 定远县| 苍溪县| 娱乐| 招远市| 沽源县| 永泰县| 大理市| 沙雅县| 柳江县| 南木林县| 阿尔山市| 思茅市| 海安县| 普陀区| 自治县| 沙田区| 台州市| 宝山区| 景宁| 论坛| 永靖县| 哈密市| 涟源市| 淄博市| 苍山县| 余姚市| 同江市| 四川省| 六盘水市| 奉贤区| 南投县| 焦作市| 上饶市| 宿州市| 南漳县| 化德县| 勃利县| 吴桥县| 木兰县|