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

如何通过Golang Pulsar API检测与Broker的连接断开?

Pulsar Go客户端检测Broker连接断开的解决方案

你遇到的问题是因为Pulsar Go客户端默认启用了自动重连机制,连接断开后会在后台持续重试,不会立刻让Receive()返回错误。以下是两种直接检测连接状态的方法:

1. 使用客户端连接状态监听器

在创建客户端时配置ConnectionListener,它会实时监听连接状态变化,包括断开、重连、正在连接等事件:

import (
    "fmt"
    "context"
    "time"

    "github.com/apache/pulsar-client-go/pulsar"
)

func main() {
    client, err := pulsar.NewClient(pulsar.ClientOptions{
        URL: "pulsar://localhost:6650",
        ConnectionListener: func(event pulsar.ConnectionStateChangedEvent) {
            switch event.State {
            case pulsar.ConnectionClosed:
                fmt.Println("[连接状态] 与Broker的连接已断开")
                // 这里可添加断开后的处理逻辑,比如告警、切换备用节点等
            case pulsar.ConnectionConnected:
                fmt.Println("[连接状态] 与Broker的连接已恢复")
            case pulsar.ConnectionConnecting:
                fmt.Println("[连接状态] 正在尝试重连Broker")
            }
        },
    })
    if err != nil {
        fmt.Println("创建客户端失败:", err)
        return
    }
    defer client.Close()

    consumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:            "persistent://public/default/my-topic",
        SubscriptionName: "my-sub",
    })
    if err != nil {
        fmt.Println("创建消费者失败:", err)
        return
    }
    defer consumer.Close()

    for {
        msg, err := consumer.Receive(context.Background())
        if err != nil {
            fmt.Println("接收消息出错:", err)
            break
        }
        fmt.Printf("收到消息: %s\n", string(msg.Payload()))
        consumer.Ack(msg)
    }
}

2. 配置消费者操作超时+错误判断

通过设置OperationTimeout让Receive()在连接异常且超时后返回错误,再结合错误类型判断是否为连接断开:

consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "persistent://public/default/my-topic",
    SubscriptionName: "my-sub",
    OperationTimeout: 30 * time.Second, // 设置操作超时时间
})

当连接断开且超过30秒无法恢复时,Receive()会返回错误,你可以在错误处理中精准判断:

msg, err := consumer.Receive(context.Background())
if err != nil {
    if err == pulsar.ErrConnectionClosed || err == pulsar.ErrTimeout {
        fmt.Println("检测到与Broker的连接丢失")
    } else {
        fmt.Println("接收消息出错:", err)
    }
    break
}

关键说明

  • Pulsar客户端默认的自动重连逻辑会让Receive()持续阻塞直到重连成功,这是你代码中未收到错误的核心原因。
  • ConnectionListener是最直接的状态感知方式,无需依赖Receive()的错误返回,能实时获取连接状态变化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:40:11