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

Pulsar如何处理Broker与Consumer间ACK确认丢失问题?Golang配置咨询

Pulsar对ACK确认丢失场景的处理机制

Broker端处理逻辑

  • Broker收到消费者的ACK请求后,首先将对应的消息偏移量持久化到BookKeeper,只有持久化操作成功完成,才会向消费者返回ACK确认。
  • 若在返回确认前连接中断,只要偏移量已经持久化,后续消费者重新连接时,会从已确认的最新偏移量开始消费,不会重复处理已ACK的消息。
  • 若偏移量尚未完成持久化连接就中断,Broker重启后会保留之前的订阅位置,消费者重新连接后会重新投递未确认的消息,保证消息不丢失。

消费者端处理逻辑

  • 消费者发送ACK后若未收到Broker的确认,默认不会自动重试ACK操作,需要开发者通过客户端API的返回结果自行处理重试逻辑。
  • Pulsar客户端会对ACK请求设置超时时间,超时后会返回错误,开发者可基于此实现有限次数的重试策略;同时Broker对重复ACK做幂等处理,多次发送同一消息的ACK不会引发异常。
Golang客户端相关配置与代码示例

关键配置项

  • OperationTimeout:设置客户端所有操作(包括ACK)的超时时间,单位为时间类型(如time.Second * 30)。当ACK操作超过该时间未收到Broker确认,会返回超时错误,触发重试逻辑。
  • MaxReconnectToBroker:设置客户端与Broker断开后的最大重连次数,保证消费者能重新恢复连接,获取最新的订阅偏移量。

代码示例

package main

import (
    "context"
    "log"
    "time"

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

func main() {
    client, err := pulsar.NewClient(pulsar.ClientOptions{
        URL: "pulsar://localhost:6650",
    })
    if err != nil {
        log.Fatalf("Failed to create client: %v", err)
    }
    defer client.Close()

    consumer, err := client.Subscribe(pulsar.ConsumerOptions{
        Topic:               "my-topic",
        SubscriptionName:    "my-subscription",
        OperationTimeout:    30 * time.Second, // ACK操作超时时间
        MaxReconnectToBroker: 5, // 最大重连次数
    })
    if err != nil {
        log.Fatalf("Failed to create consumer: %v", err)
    }
    defer consumer.Close()

    for {
        msg, err := consumer.Receive(context.Background())
        if err != nil {
            log.Printf("Failed to receive message: %v", err)
            continue
        }

        // 处理消息逻辑
        log.Printf("Received message: %s", string(msg.Payload()))

        // 发送ACK并等待结果
        ackFuture := consumer.Ack(msg)
        err = ackFuture.Complete(context.Background())
        if err != nil {
            log.Printf("ACK failed for message %s: %v, retrying...", msg.MessageID(), err)
            // 实现有限次数重试
            retryCount := 0
            maxRetries := 3
            for retryCount < maxRetries {
                time.Sleep(2 * time.Second)
                ackFuture := consumer.Ack(msg)
                err = ackFuture.Complete(context.Background())
                if err == nil {
                    log.Printf("ACK succeeded after retry %d", retryCount+1)
                    break
                }
                retryCount++
                log.Printf("Retry %d failed: %v", retryCount+1, err)
            }
            if retryCount == maxRetries {
                log.Printf("Failed to ACK message %s after %d retries", msg.MessageID(), maxRetries)
                // 可根据业务需求处理,比如死信队列
            }
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 10:25:30