You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Golang中ZMQ PUB/SUB订阅者无法接收全部消息的问题排查

ZMQ PUB/SUB模式消息丢失与异常延迟问题分析

问题背景

使用Golang基于ZMQ实现PUB/SUB通信,尝试从PUB端发送指定数量(如10000条)消息到SUB端,测试接收耗时。但SUB端始终无法接收全部消息,且某次运行时在接收3000条后出现异常延迟。

PUB端代码

package main

import (
    "fmt"
    "log"
    "os"
    "strconv"
    "time"

    zmq "github.com/pebbe/zmq4"
)

const letters = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ"

func SampleStrOfSize(size int) string {
    b := make([]byte, size)
    idx := 0
    for i := 0; i < size; i++ {
        b[i] = letters[idx]
        idx = (idx + 1) % len(letters)
    }
    // fmt.Println("message: ", string(b))
    return string(b)
}

func main() {
    publisher, err := zmq.NewSocket(zmq.PUB)
    if err != nil {
        os.Exit(1)
    }
    defer publisher.Close()
    connectionStr := "tcp://127.0.0.1:5555"

    err = publisher.Bind(connectionStr)
    if err != nil {
        log.Fatal(err)
    }
    fmt.Println("Bound to", connectionStr)

    // 生成固定长度的测试消息
    message := SampleStrOfSize(1024)
    msgCount := 0

    if len(os.Args) > 1 {
        if msgCount, err = strconv.Atoi(os.Args[1]); err != nil {
            return
        }
    } else {
        fmt.Println("No msg count provided")
        return
    }

    if msgCount <= 0 {
        fmt.Printf("Invalid msg count (%v) provided", msgCount)
        return
    }

    // 等待SUB端连接
    time.Sleep(time.Second * 5)
    fmt.Println("Starting message sending")
    start := time.Now()

    for i := 0; i < msgCount; i++ {
        // 循环发送相同消息
        _, err = publisher.Send(message, 0)
        if err != nil {
            fmt.Printf("Error occured while sending message. %v", err)
        }
    }

    elapsed := time.Since(start)
    fmt.Printf("Sending %d messages took %s\n", msgCount, elapsed)
    // 延迟关闭,避免提前断开连接
    time.Sleep(time.Second * 30)
}

SUB端代码

package main

import (
    "fmt"
    "os"
    "strconv"
    "time"

    zmq "github.com/pebbe/zmq4"
)

func main() {
    subscriber, err := zmq.NewSocket(zmq.SUB)
    if err != nil {
        fmt.Println("Failed to open socket")
        os.Exit(1)
    }
    defer subscriber.Close()

    err = subscriber.Connect("tcp://127.0.0.1:5555")
    if err != nil {
        fmt.Println("Connect failed")
        os.Exit(1)
    }

    msgCount := 0
    count := 0
    if len(os.Args) > 1 {
        if msgCount, err = strconv.Atoi(os.Args[1]); err != nil {
            return
        }
    } else {
        fmt.Println("No msg count provided")
        os.Exit(1)
    }
    fmt.Printf("Expecting %d messages\n", msgCount)

    // 订阅所有消息
    err = subscriber.SetSubscribe("")
    if err != nil {
        fmt.Println("Failed to subscribe for all messages")
        os.Exit(1)
    }

    var start time.Time
    for {
        _, err := subscriber.Recv(0)

        if count == 0 {
            start = time.Now()
        }
        if err != nil {
            fmt.Println("Receive failed")
        }

        count++

        if count == msgCount {
            break
        } else if 0 == count%1000 {
            // 每接收1000条消息打印耗时
            elapsed := time.Since(start)
            fmt.Printf("Received %d messages in %s\n", count, elapsed)
        }
    }

    elapsed := time.Since(start)
    fmt.Printf("Received %d messages in %s\n", msgCount, elapsed)
}

异常运行输出

Expecting 10000 messages
Received 1000 messages in 17.435321ms
Received 2000 messages in 25.530057ms
Received 3000 messages in 27.80558ms
Received 4000 messages in 1m40.583143061s
Received 5000 messages in 1m40.590513201s
Received 6000 messages in 1m40.597145666s

问题分析与解决方案

1. 消息丢失的原因及解决方法

原因

ZMQ的PUB/SUB模式本质是无可靠投递保障的广播模式,核心问题在于:

  • 默认情况下,PUB和SUB套接字都有消息队列高水位线(ZMQ_SNDHWM/ZMQ_RCVHWM),当PUB端发送速度远快于SUB端处理速度时,PUB的发送队列会被迅速填满,后续消息会直接被丢弃。
  • 即使设置了5秒等待SUB连接,ZMQ不会主动同步连接就绪状态,极端情况下仍可能有少量消息在SUB完全订阅前被发送,但本次场景主要是队列溢出导致的大量丢消息。

解决方法

  • 调整高水位线:将PUB和SUB的高水位线设置为0(表示无队列长度限制,需注意内存占用),避免队列溢出丢消息。在创建套接字后添加:
    // PUB端设置发送高水位线
    publisher.SetSndhwm(0)
    // SUB端设置接收高水位线
    subscriber.SetRcvHwm(0)
    
  • 增加流量控制:如果是单SUB场景,可以在PUB/SUB基础上增加反向确认通道,SUB每接收一定数量的消息就向PUB发送确认信号,PUB收到确认后再继续发送,确保SUB处理完一批再发下一批。
  • 优化SUB接收效率:当前SUB采用单线程阻塞接收,可改用批量接收(RecvMessage)或多线程处理,提升消息消费速度,减少队列积压。

2. 接收3000条后异常延迟的原因

延迟出现的核心原因是消息队列溢出触发的阻塞与TCP流量控制:

  • 当PUB的发送队列被填满后,由于代码中Send使用了0(阻塞模式),PUB端会暂停发送,直到SUB端接收消息释放队列空间。
  • SUB端前3000条消息是从ZMQ的本地接收队列快速读取,当队列空了之后,需要等待PUB端重新发送,此时TCP缓冲区也可能被填满,触发TCP滑动窗口机制,导致收发双方进入等待状态,延迟陡增。

内容的提问来源于stack exchange,提问作者SomeNoob

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 23:32:05