C#8 Lambda向RabbitMQ发消息偶现AlreadyClosedException问题排查
我有一个用C# 8编写的Lambda,用来向RabbitMQ队列发送消息。功能大部分时候正常,但偶尔会记录如下异常。Lambda执行时长通常很短(<1秒)。
我用AddSingleton()注册IChannel,这是不是有问题?
原依赖注入配置
.AddSingleton<IChannel>(provider => { var rabbitMqServiceOptions = provider.GetService<IOptions<RabbitMqServiceOptions>>()!.Value; var connectionFactory = new ConnectionFactory { Uri = new Uri(rabbitMqServiceOptions.Uri) }; IConnection conn = connectionFactory.CreateConnectionAsync().GetAwaiter().GetResult(); IChannel channel = conn.CreateChannelAsync().GetAwaiter().GetResult(); return channel; }) .AddTransient<RabbitMqServiceOptions>(provider => { return provider.GetService<IOptions<RabbitMqServiceOptions>>()!.Value; })
异常信息
Exception:RabbitMQ.Client.Exceptions.AlreadyClosedException: Already closed: The AMQP operation was interrupted: AMQP close-reason, initiated by Library, code=0, text='End of stream', classId=0, methodId=0, exception=System.IO.EndOfStreamException: Pipe is completed. at RabbitMQ.Client.Impl.InboundFrame.ReadFromPipeAsync(PipeReader reader, UInt32 maxInboundMessageBodySize, InboundFrame frame, CancellationToken mainLoopCancellationToken) at RabbitMQ.Client.Framing.Connection.ReceiveLoopAsync(CancellationToken mainLoopCancellationToken) at RabbitMQ.Client.Framing.Connection.MainLoop() ---> System.IO.EndOfStreamException: Pipe is completed. at RabbitMQ.Client.Impl.InboundFrame.ReadFromPipeAsync(PipeReader reader, UInt32 maxInboundMessageBodySize, InboundFrame frame, CancellationToken mainLoopCancellationToken) at RabbitMQ.Client.Framing.Connection.ReceiveLoopAsync(CancellationToken mainLoopCancellationToken) at RabbitMQ.Client.Framing.Connection.MainLoop() --- End of inner exception stack trace --- at RabbitMQ.Client.Impl.SessionBase.ThrowAlreadyClosedException() at RabbitMQ.Client.Impl.SessionBase.TransmitAsync[T](T& cmd, CancellationToken cancellationToken) at RabbitMQ.Client.Impl.Channel.QueueDeclareAsync(String queue, Boolean durable, Boolean exclusive, Boolean autoDelete, IDictionary
2 arguments, Boolean passive, Boolean noWait, CancellationToken cancellationToken) at RabbitMQ.Client.Impl.AutorecoveringChannel.QueueDeclareAsync(String queue, Boolean durable, Boolean exclusive, Boolean autoDelete, IDictionary2 arguments, Boolean passive, Boolean noWait, CancellationToken cancellationToken)
编辑后的配置
我已经更新了注册配置,把IChannel改成Transient,希望能解决问题:
.AddSingleton<IConnection>(provider => { var rabbitMqServiceOptions = provider.GetService<IOptions<RabbitMqServiceOptions>>()!.Value; var connectionFactory = new ConnectionFactory { Uri = new Uri(rabbitMqServiceOptions.Uri) }; IConnection conn = connectionFactory.CreateConnectionAsync().GetAwaiter().GetResult(); return conn; }) .AddTransient<IChannel>(provider => { var conn = provider.GetService<IConnection>(); IChannel channel = conn!.CreateChannelAsync().GetAwaiter().GetResult(); return channel; }) .AddTransient<RabbitMqServiceOptions>(provider => { return provider.GetService<IOptions<RabbitMqServiceOptions>>()!.Value; })
回答
你最初用AddSingleton注册IChannel确实是问题根源。
RabbitMQ的Channel是和Connection绑定的轻量级会话,它不是线程安全的,也不适合作为单例长期持有。Lambda是无服务器环境,执行完成后会被冻结,下次唤醒时,之前创建的单例Channel可能已经因为连接超时、RabbitMQ端主动关闭等原因失效,这就会抛出AlreadyClosedException。
你更新后的配置是正确的方向:
- 保持
IConnection为单例:RabbitMQ的Connection线程安全,创建成本高,复用Connection是最佳实践。 - 将
IChannel改为Transient:每次需要发送消息时创建新的Channel,使用完毕后及时关闭(可在服务类中用using包裹,或依赖DI的Transient生命周期自动回收),避免复用失效的Channel。
额外建议:
- 避免用
.GetAwaiter().GetResult()同步阻塞异步方法,.NET 6+支持AddSingletonAsync等异步DI注册方式,能避免阻塞带来的性能问题或死锁风险。 - 推荐在使用Channel时用
using语句确保及时释放:using var channel = _connection.CreateChannel(); // 执行消息发送逻辑
内容的提问来源于stack exchange,提问作者deejbee

