Channel读取正常运行一段时间后停滞问题求助
我有一个长期运行的类(非BackgroundService),内部使用Channel缓冲传入请求。类中包含Enqueue方法和消息读取循环_MainMessageLoop,代码如下:
Enqueue方法代码
public async Task<bool> Enqueue(IMessage message) { bool result; try { if (_cts.IsCancellationRequested) { _logger.LogInformation("{TypeName} cannot be enqueued because the queue is stopped.", message.GetType().Name); result = false; } else { await _messageChannel.Writer.WriteAsync(message); result = true; _logger.LogInformation("Message {TypeName} written to internal channel.", message.GetType().Name); } } catch (ChannelClosedException) { _logger.LogInformation("{TypeName} cannot be enqueued because the channel is closed.", message.GetType().Name); result = false; } catch (Exception e) { _logger.LogError(e, "Error enqueuing {TypeName}", message.GetType().Name); result = false; } return result; }
消息读取循环代码
在类的构造函数中通过Task.Run(() => _MainMessageLoop(_cts.Token))启动循环:
private async Task _MainMessageLoop(CancellationToken stoppingToken) { try { var channelReader = _messageChannel.Reader; while (await channelReader.WaitToReadAsync(stoppingToken).ConfigureAwait(false)) { try { var item = await channelReader.ReadAsync(stoppingToken).ConfigureAwait(false); _logger.LogInformation("Message {TypeName} read from internal channel.", item.GetType().Name); await _DoSomethingAsync(item); } catch (Exception e) { _logger.LogError(e, "Error doing something"); } } _logger.LogInformation("Channel reader loop exited {MethodName}.", nameof(_MainMessageLoop)); } catch (OperationCanceledException) { _logger.LogDebug("OperationCanceled {MethodName}", nameof(_MainMessageLoop)); throw; } catch (Exception e) { _logger.LogError(e, "Error in {MethodName}", nameof(_MainMessageLoop)); throw; } }
问题现象
程序在极低消息负载下能正常运行数小时,但之后读取端会停滞。日志显示写入操作仍在进行,但读取循环无报错、无异常日志,只是停止工作。尝试过Bounded/Unbounded Channel、TryRead与ReadAsync的多种组合,均未解决问题。应用中其他HostedService级别的单例Channel运行正常。
1. _DoSomethingAsync阻塞或死锁
这是最常见的触发因素:如果_DoSomethingAsync内部存在同步阻塞(比如调用.Result、.Wait())或死锁,会导致读取循环的异步任务被卡住,无法继续处理后续消息。即使外部日志显示写入正常,读取端会停留在await _DoSomethingAsync(item)步骤。
验证:在_DoSomethingAsync的入口和出口添加日志,观察是否有消息进入后无出口日志输出。
修复:
- 确保
_DoSomethingAsync全程使用异步模式,避免任何同步阻塞操作; - 若必须调用同步代码,用
Task.Run包裹以避免占用读取循环线程:await Task.Run(() => _DoSomethingSync(item), stoppingToken);
2. Task.Run启动的任务未被跟踪
非BackgroundService类中,Task.Run返回的任务若未被持续引用,可能被GC视为可回收对象,导致任务被静默终止。而HostedService的任务由框架跟踪,不会出现此问题。
修复:
- 在类中添加私有字段保存任务实例:
private Task _messageLoopTask; - 构造函数中初始化并保存:
_messageLoopTask = Task.Run(() => _MainMessageLoop(_cts.Token), _cts.Token); - 实现
IDisposable接口,在销毁逻辑中等待任务完成:public void Dispose() { _cts.Cancel(); _messageLoopTask?.Wait(); _cts.Dispose(); _messageChannel.Writer.Complete(); }
3. CancellationToken被意外触发
检查_cts的创建和调用逻辑,是否存在意外调用Cancel()的情况。若stoppingToken被取消,WaitToReadAsync会返回false,循环退出,此时读取端停止工作,但写入端仅在每次Enqueue时检查取消状态,若取消操作发生在两次Enqueue之间,后续写入才会触发失败日志。
验证:将_MainMessageLoop中循环退出的日志级别升级为Warning,观察是否有该日志输出。
修复:
- 确保
_cts仅在类销毁逻辑中调用Cancel(); - 在
Enqueue方法中,将_cts.Token传入WriteAsync,让写入操作能及时响应取消:await _messageChannel.Writer.WriteAsync(message, _cts.Token);
4. Channel Writer被意外完成
若_messageChannel.Writer.Complete()被意外调用,WaitToReadAsync会返回false导致循环退出。但此时写入端调用WriteAsync会抛出ChannelClosedException,若用户日志中无相关记录,此可能性较低,但仍需排查所有调用Complete()的代码路径,确保仅在类销毁时执行。
内容的提问来源于stack exchange,提问作者micahtan

