Go语言中RabbitMQ通道关闭后消费者未收到通知的问题
为什么将RabbitMQ消费通道移到单独goroutine后才能收到通道关闭通知?
我用Go开发,RabbitMQ版本3.12.1,服务器因错误「PRECONDITION_FAILED - delivery acknowledgement on channel 1 timed out. Timeout value used: 1000 ms.」关闭通道,但原代码里的消费者始终收不到「channel is closed」的关闭通知。
原代码将consumeChan、chanClosedChan、connClosedChan和ctx.Done()放在同一个select中监听:
amqpConn, amqpChannel, err := connectToAmqp(ctx, config, logger) if err != nil { return nil, err } consumeChan, err := amqpChannel.Consume("TestQueue", "", false, false, false, false, nil) chanClosedChan := amqpChannel.NotifyClose(makeAmqpErrorChan()) connClosedChan := amqpConn.NotifyClose(makeAmqpErrorChan()) for { select { case amqpErr := <-chanClosedChan: logger.Error(ctx, "channel is closed") case amqpErr := <-connClosedChan: logger.Error(ctx, "closed rabbitmq connection") case <-ctx.Done(): logger.Info(ctx, "Closing rabbitmq connection") err := amqpChannel.Close() if err != nil { logger.Error(ctx, "Failed to cleanly close rabbitmq connection", field.Error(err)) } return case amqpMessage, ok := <-consumeChan: if ok { logger.Debug(ctx, "Received amqpMessage message", field.Any("AmqpMessage", amqpMessage)) } } }
把consumeChan的消费逻辑移到单独goroutine后,就能正常收到通道关闭的通知了:
// ... 省略连接代码 go handleConsumerChannel(ctx, consumeChan) for { select { case amqpErr := <-chanClosedChan: logger.Error(ctx, "channel is closed") case amqpErr := <-connClosedChan: logger.Error(ctx, "closed rabbitmq connection") case <-ctx.Done(): logger.Info(ctx, "Closing rabbitmq connection") err := amqpChannel.Close() if err != nil { logger.Error(ctx, "Failed to cleanly close rabbitmq connection", field.Error(err)) } return } } // ... func handleConsumerChannel(ctx context.Context, consumeChan <-chan amqp.Delivery) { for { select { case amqpMessage, ok := <-consumeChan: if ok { logger.Debug(ctx, "Received amqpMessage message", field.Any("AmqpMessage", amqpMessage)) } } } }
问题原因
核心在于Go语言select语句的执行机制,以及RabbitMQ的consumeChan在通道关闭后的行为:
- 当RabbitMQ服务器关闭通道时,
amqpChannel会关闭对应的consumeChan。此时对consumeChan的读操作会立即返回(零值, false),这个case会持续处于就绪状态。 - 在原代码的
select中,由于consumeChan的case始终就绪,Go的select会反复选中这个分支(即使ok为false,代码也没有退出循环或终止该分支的处理),导致chanClosedChan的消息完全没有机会被处理,自然看不到通道关闭的日志。 - 修改后,消费逻辑被移到单独goroutine,主goroutine的
select不再包含这个始终就绪的case。当chanClosedChan有消息时,主goroutine的select能正常选中该分支,处理通道关闭的通知。
另外补充:原代码中consumeChan分支在ok为false时没有任何退出逻辑,会导致该goroutine无限循环处理关闭通道的读返回,进一步占用了所有select的执行机会。
内容的提问来源于stack exchange,提问作者Aryaman Gupta
相关产品推荐
相关产品推荐

