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

使用HotChocolate v12订阅时的内存泄漏排查求助

HotChocolate v12.11.1订阅功能内存泄漏排查

结合你提供的代码,以下是可能导致内存泄漏的核心问题及修复建议:

1. RxSouceStream未实现资源清理逻辑

你的RxSouceStream类的DisposeAsync方法直接返回default,完全没有处理底层订阅资源的释放:

  • 扩展方法中通过Task.WhenAll获取到的ISourceStream<TMessage>实例(每个对应一个股票标的的订阅),大概率实现了IAsyncDisposable接口,但当前代码既没有持有这些实例的引用,也没有在DisposeAsync中释放它们。
  • 当客户端断开订阅时,HotChocolate会调用RxSouceStream.DisposeAsync,但此时底层的topic订阅仍然保持活跃,导致相关对象无法被GC回收,引发内存泄漏。

修复建议:
修改RxSouceStream,持有底层ISourceStream实例并在DisposeAsync中释放:

public class RxSouceStream<TSource> : ISourceStream<TSource>, IAsyncDisposable
{
    private readonly IAsyncEnumerable<TSource> _enumerable;
    private readonly ISourceStream<TSource>[] _underlyingStreams;

    public RxSouceStream(IAsyncEnumerable<TSource>[] sources, ISourceStream<TSource>[] underlyingStreams)
    {
        _enumerable = AsyncEnumerableEx.Merge(sources);
        _underlyingStreams = underlyingStreams;
    }
    
    public async ValueTask DisposeAsync()
    {
        // 遍历释放所有底层订阅流
        foreach (var stream in _underlyingStreams)
        {
            if (stream is IAsyncDisposable asyncDisposable)
            {
                await asyncDisposable.DisposeAsync().ConfigureAwait(false);
            }
        }
    }

    public IAsyncEnumerable<TSource> ReadEventsAsync()
    {
        return _enumerable;
    }

    IAsyncEnumerable<object> ISourceStream.ReadEventsAsync()
    {
        return _enumerable.Select(x => (object)x!);
    }
}

同时修改扩展方法,将底层流实例传递给RxSouceStream:

public static async ValueTask<ISourceStream<TMessage>> SubscribeAsync<TTopic, TMessage>(
    this ITopicEventReceiver topicEventReceiver, 
    IEnumerable<TTopic> topics, 
    CancellationToken cancellationToken = default)
    where TTopic : notnull
{
    var subscriptions = topics.Take(100)
        .Select(p => topicEventReceiver.SubscribeAsync<TTopic, TMessage>(p, cancellationToken).AsTask())
        .ToArray();

    // 保存底层订阅流实例
    var underlyingStreams = await Task.WhenAll(subscriptions);
    var streams = underlyingStreams.Select(p => p.ReadEventsAsync()).ToArray();

    return new RxSouceStream<TMessage>(streams, underlyingStreams);
}

2. 合并流未关联取消令牌

当前合并后的流没有绑定传入的CancellationToken,当客户端断开(HotChocolate会触发令牌取消)时,合并流无法及时终止,导致底层订阅继续运行并占用内存。

修复建议:
在合并流时应用取消令牌:

public RxSouceStream(IAsyncEnumerable<TSource>[] sources, ISourceStream<TSource>[] underlyingStreams, CancellationToken cancellationToken)
{
    _enumerable = AsyncEnumerableEx.Merge(sources).WithCancellation(cancellationToken);
    _underlyingStreams = underlyingStreams;
}

同时在扩展方法中传递令牌:

return new RxSouceStream<TMessage>(streams, underlyingStreams, cancellationToken);

3. 潜在的流重复订阅风险

ReadEventsAsync方法每次都返回同一个合并后的_enumerable实例,如果HotChocolate多次调用该方法(比如订阅重试场景),可能导致重复订阅底层topic,进而累积未释放的资源。

修复建议:
使用AsyncEnumerable.Cache()缓存流,避免重复枚举:

public RxSouceStream(IAsyncEnumerable<TSource>[] sources, ISourceStream<TSource>[] underlyingStreams, CancellationToken cancellationToken)
{
    _enumerable = AsyncEnumerableEx.Merge(sources)
        .WithCancellation(cancellationToken)
        .Cache(); // 缓存流,避免重复枚举触发多次订阅
    _underlyingStreams = underlyingStreams;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 18:32:12