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
相关产品推荐
相关产品推荐

