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

