如何通过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
相关产品推荐
相关产品推荐

