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

AMQP Golang生产者发送大消息时丢失最后一条消息的问题

问题:大消息发布时最后一条丢失,仅添加Sleep后正常

我需要从文件读取每行一条的消息发布到RabbitMQ(使用mTLS),大部分情况运行正常,但遇到几十KB的大消息时,最后一条完全无法发布。

生产者代码

func main() {
        flag.Parse()

        cfg := new(tls.Config)
        cfg.RootCAs = x509.NewCertPool()

        caCert, err := os.ReadFile(*caFile)
        failOnError(err, "Unable to read CA bundle")
        cfg.RootCAs.AppendCertsFromPEM(caCert)

        cert, err := tls.LoadX509KeyPair(*certFile, *keyFile)
        failOnError(err, "Unable to read certificate or key")
        cfg.Certificates = append(cfg.Certificates, cert)

        conn, err := amqp.DialTLS(*url, cfg)
        failOnError(err, "Unable to dial with TLS")
        defer conn.Close()

        ch, err := conn.Channel()
        failOnError(err, "Failed to open a channel")

        messages, err := getMessages()
        failOnError(err, "Unable to get messages from file")
        log.Printf("%d messages to publish", len(messages))

        for i, m := range messages {
                log.Printf("publish message, idx: %d, len: %d", i, len(m))
                err = ch.Publish(
                        *exchange,   // exchange
                        *routingKey, // routing key
                        false,
                        false, // immediate
                        amqp.Publishing{
                                ContentType:  "text/plain",
                                Body:         []byte(m),
                                DeliveryMode: amqp.Persistent,
                        })
                failOnError(err, "Failed to publish a message")

        }
        time.Sleep(0 * time.Millisecond) // special reason explained at bottom
}

运行日志

生产者日志

$ go run producer-file.go --cert client.cer --key client.key --ca ca-bundle.pem --url amqps://login:password@server:5671/virtualhost --exchange test_exchange --routing-key test_routing_key --file test_03_bundle4.txt
2023/03/21 12:01:02 4 messages to publish
2023/03/21 12:01:02 publish message, idx: 0, len: 169403
2023/03/21 12:01:02 publish message, idx: 1, len: 169148
2023/03/21 12:01:02 publish message, idx: 2, len: 169659
2023/03/21 12:01:02 publish message, idx: 3, len: 86000

消费者日志

$ go run consumer.go --cert client.cer --key client.key --ca ca-bundle.pem --url amqps://login:password@server:5671/virtualhost -p 10 --queue test_queue
2023/03/21 12:01:02 Received a message, len: 169403
2023/03/21 12:01:02 Received a message, len: 169148
2023/03/21 12:01:02 Received a message, len: 169659
2023/03/21 12:01:02 Current speed: 3 msg/s

现象与疑问

日志显示4条消息都执行了发布操作,但消费者仅收到前3条。在循环后添加20ms的time.Sleep()后,所有消息都能被正常接收。短消息(约20字节)无需Sleep即可正常发布。

明明defer conn.Close()应该确保缓冲区排空前不会关闭连接,为什么大消息的最后一条会丢失?


原因与解决方案

核心原因

amqp.Publish是异步非阻塞操作:它只是将消息写入本地缓冲区,并没有等待消息真正发送到RabbitMQ Broker,也不会等待Broker的确认。

defer conn.Close()虽然会关闭连接,但关闭逻辑并不会同步等待所有缓冲区中的消息完成传输并被Broker确认:

  • 短消息体积小,在程序退出前,网络栈已经完成发送,Broker也处理完毕,因此不会丢失。
  • 大消息需要更多时间完成网络传输,加上你设置了DeliveryMode: amqp.Persistent,Broker需要将消息写入磁盘做持久化,耗时更长。程序发布最后一条消息后立刻触发conn.Close(),此时消息可能还在本地缓冲区,或刚发出但未被Broker确认,连接断开直接导致消息丢失。

Sleep之所以有用,是因为给了足够时间让最后一条消息完成传输和Broker处理,避免连接提前关闭,但这种方式不可靠——不同环境需要的等待时间差异很大,无法保证通用性。

正确解决方案:使用发布确认机制

开启RabbitMQ的发布确认模式,确保所有消息都被Broker确认接收后再退出程序:

func main() {
        // ... 原有初始化代码(TLS配置、连接、信道创建)

        messages, err := getMessages()
        failOnError(err, "Unable to get messages from file")
        log.Printf("%d messages to publish", len(messages))

        // 1. 开启信道的发布确认模式
        ch.Confirm(false)
        // 2. 创建确认通道,用于接收Broker的确认信号
        confirmChan := ch.NotifyPublish(make(chan amqp.Confirmation, len(messages)))

        for i, m := range messages {
                log.Printf("publish message, idx: %d, len: %d", i, len(m))
                err = ch.Publish(
                        *exchange,   // exchange
                        *routingKey, // routing key
                        false,
                        false, // immediate
                        amqp.Publishing{
                                ContentType:  "text/plain",
                                Body:         []byte(m),
                                DeliveryMode: amqp.Persistent,
                        })
                failOnError(err, "Failed to publish a message")
        }

        // 3. 等待所有消息被Broker确认
        for i := 0; i < len(messages); i++ {
                confirm := <-confirmChan
                if !confirm.Ack {
                        log.Fatalf("Message %d not acknowledged by broker", i)
                }
        }

        // 此时所有消息都已确认,无需Sleep,直接退出即可
}

这种方式从根本上保证了消息的可靠性,不需要依赖Sleep这种不可靠的临时方案。


内容的提问来源于stack exchange,提问作者Mariusz Jędrzejewski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 03:48:07