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

如何解决NatsSubscriber中Disposable.Create无法调用异步DisposeAsync的问题

解决NatsSubscriber中异步释放导致的async void问题

核心问题

Disposable.Create的回调是同步委托,直接在里面使用await connection.DisposeAsync()会生成async void方法,无法捕获异常且无法等待异步操作完成,存在资源泄漏风险。

解决方案:自定义异步处置类

实现同时兼容IDisposable和IAsyncDisposable的自定义处置类,替代Disposable.Create,确保异步释放逻辑在正确的异步上下文中执行。

步骤1:自定义处置类

创建专门处理NATS订阅释放的处置类,包含异步释放逻辑:

private sealed class NatsSubscriptionDisposable : IDisposable, IAsyncDisposable
{
    private readonly NatsConnection _connection;
    private readonly IObserver<TMessage> _observer;
    private readonly string _subject;
    private bool _disposed;

    public NatsSubscriptionDisposable(NatsConnection connection, IObserver<TMessage> observer, string subject)
    {
        _connection = connection;
        _observer = observer;
        _subject = subject;
    }

    // 同步Dispose:内部转异步并阻塞,兼容旧调用方式
    public void Dispose()
    {
        DisposeAsync().AsTask().GetAwaiter().GetResult();
    }

    // 异步释放:正确处理异步操作
    public async ValueTask DisposeAsync()
    {
        if (_disposed) return;
        _disposed = true;

        // 1. 取消NATS订阅(如果AlterNats提供UnsubscribeAsync方法)
        await _connection.UnsubscribeAsync(_subject).ConfigureAwait(false);
        // 2. 通知观察者完成
        _observer.OnCompleted();
        // 3. 释放NATS连接
        await _connection.DisposeAsync().ConfigureAwait(false);
    }
}

步骤2:修改SubscribeAsync方法

替换原有的Disposable.Create,返回自定义的处置类:

public async Task<IDisposable> SubscribeAsync(TKey key, Action<TMessage> callback)
{
    var subject = GetSubjectString(key);

    if (subject == null)
    {
        throw new ArgumentNullException(nameof(key));
    }

    var connection = await _connectionFactory.GetConnectionAsync().ConfigureAwait(false);

    var observer = Observer.Create(callback);

    await connection.SubscribeAsync<TMessage>(subject, data =>
    {
        observer.OnNext(data);
    }).ConfigureAwait(false);

    // 返回自定义异步处置类
    return new NatsSubscriptionDisposable(connection, observer, subject);
}

注意事项

  • 优先使用异步释放:推荐调用者使用await using语法来异步释放订阅,避免同步Dispose带来的阻塞:
    await using var subscription = await natsSubscriber.SubscribeAsync(key, msg => 
    {
        // 消息处理逻辑
    });
    
  • 共享连接场景调整:如果NatsConnectionFactory返回的是共享连接(多个订阅复用同一连接),不要直接释放连接,而是保存SubscribeAsync返回的订阅对象(如ISubscription),调用其UnsubscribeAsync方法取消订阅即可,避免影响其他订阅:
    // 假设SubscribeAsync返回ISubscription
    var subscription = await connection.SubscribeAsync<TMessage>(subject, data =>
    {
        observer.OnNext(data);
    }).ConfigureAwait(false);
    
    // 调整处置类保存subscription而非connection
    private sealed class NatsSubscriptionDisposable : IDisposable, IAsyncDisposable
    {
        private readonly ISubscription _subscription;
        private readonly IObserver<TMessage> _observer;
        private bool _disposed;
    
        public NatsSubscriptionDisposable(ISubscription subscription, IObserver<TMessage> observer)
        {
            _subscription = subscription;
            _observer = observer;
        }
    
        public void Dispose()
        {
            DisposeAsync().AsTask().GetAwaiter().GetResult();
        }
    
        public async ValueTask DisposeAsync()
        {
            if (_disposed) return;
            _disposed = true;
    
            await _subscription.UnsubscribeAsync().ConfigureAwait(false);
            _observer.OnCompleted();
        }
    }
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 19:20:47