并行生产者-消费者模式下,如何让消费者感知生产者异常?
解决生产者异常时消费者感知问题的几种可行方案
方案一:追踪生产者Task,捕获异常并统一通知
把生产者任务保存下来,在生产者内部捕获异常,不管正常完成还是抛出异常,都标记完成状态,同时记录异常信息让消费者处理。
修改后代码示例:
public override async Task SendMultiSkuSpmRequests(SpmRequestMultiSku hostRequest, IServerStreamWriter<SpmResponse> responseStream, ServerCallContext context) { ConcurrentQueue<SpmResponse> queue = new ConcurrentQueue<SpmResponse>(); ManualResetEvent producerFinished = new ManualResetEvent(false); Exception producerException = null; // 保存生产者任务,便于后续追踪状态 var producerTask = Task.Factory.StartNew(() => { try { EnqueueSpmResponsesForMultiSku(hostRequest, queue, producerFinished); } catch (Exception ex) { producerException = ex; // 异常时必须标记完成,避免消费者永久阻塞 producerFinished.Set(); logger.Error("生产者执行异常", ex); } }); logger.Debug("开始监听Spm响应队列,将逐个通过流发送回测试主机。"); while (true) { SpmResponse response; // 优化原逻辑:直接用TryDequeue循环,避免先判断Any再Dequeue的冗余 while (queue.TryDequeue(out response)) { logger.Debug("向客户端发送Spm响应。"); await responseStream.WriteAsync(response); } // 检查生产者是否完成(正常或异常) if (producerFinished.WaitOne(50)) { // 处理生产者异常 if (producerException != null) { logger.Error("生产者异常导致流程终止", producerException); // 可选:向客户端返回错误响应,或者抛出异常终止当前调用 throw producerException; } logger.Debug("生产者已完成且队列为空 - 将退出循环"); break; } // 额外处理客户端主动断开的场景 if (context.CancellationToken.IsCancellationRequested) { logger.Debug("客户端取消请求,终止处理"); break; } } // 等待生产者任务完成,避免遗漏未捕获的异常 try { await producerTask; } catch (Exception ex) { logger.Error("生产者任务最终异常", ex); } }
方案二:用自定义状态类封装完成状态与异常信息
创建一个统一的状态类,把完成标记、异常信息打包在一起,生产者无论成功失败都更新状态,消费者循环中检查状态即可。
代码示例:
// 自定义生产者状态类,封装完成状态和异常信息 private class ProducerState { public bool IsFinished { get; set; } public Exception Exception { get; set; } public readonly ManualResetEvent FinishedEvent = new ManualResetEvent(false); } public override async Task SendMultiSkuSpmRequests(SpmRequestMultiSku hostRequest, IServerStreamWriter<SpmResponse> responseStream, ServerCallContext context) { ConcurrentQueue<SpmResponse> queue = new ConcurrentQueue<SpmResponse>(); var producerState = new ProducerState(); Task.Factory.StartNew(() => { try { EnqueueSpmResponsesForMultiSku(hostRequest, queue, producerState.FinishedEvent); producerState.IsFinished = true; } catch (Exception ex) { producerState.IsFinished = true; producerState.Exception = ex; producerState.FinishedEvent.Set(); logger.Error("生产者执行异常", ex); } }); logger.Debug("开始监听Spm响应队列,将逐个通过流发送回测试主机。"); while (true) { SpmResponse response; while (queue.TryDequeue(out response)) { logger.Debug("向客户端发送Spm响应。"); await responseStream.WriteAsync(response); } if (producerState.FinishedEvent.WaitOne(50)) { if (producerState.Exception != null) { logger.Error("生产者异常导致流程终止", producerState.Exception); throw producerState.Exception; } logger.Debug("生产者已完成且队列为空 - 将退出循环"); break; } if (context.CancellationToken.IsCancellationRequested) { logger.Debug("客户端取消请求,终止处理"); break; } } }
方案三:使用TPL Dataflow简化生产者消费者流程
TPL Dataflow自带完成、异常传播和取消机制,无需手动管理队列和事件,代码更简洁健壮。
代码示例:
using System.Threading.Tasks.Dataflow; public override async Task SendMultiSkuSpmRequests(SpmRequestMultiSku hostRequest, IServerStreamWriter<SpmResponse> responseStream, ServerCallContext context) { // 创建处理块:收到数据就发送给客户端 var sendBlock = new ActionBlock<SpmResponse>(async response => { logger.Debug("向客户端发送Spm响应。"); await responseStream.WriteAsync(response); }, new ExecutionDataflowBlockOptions { CancellationToken = context.CancellationToken }); // 生产者任务:生成数据并发送到处理块 var producerTask = Task.Run(async () => { try { ConcurrentQueue<SpmResponse> queue = new ConcurrentQueue<SpmResponse>(); EnqueueSpmResponsesForMultiSku(hostRequest, queue, null); // 不再需要手动事件 while (queue.TryDequeue(out var response)) { await sendBlock.SendAsync(response, context.CancellationToken); } } catch (Exception ex) { logger.Error("生产者执行异常", ex); // 标记处理块为故障状态,让后续感知异常 sendBlock.Fault(ex); } finally { // 标记生产者完成,处理块不再接收新数据 sendBlock.Complete(); } }, context.CancellationToken); logger.Debug("开始监听Spm响应,将逐个通过流发送回测试主机。"); try { // 等待处理块完成所有数据发送 await sendBlock.Completion; // 等待生产者任务完成,捕获可能的异常 await producerTask; } catch (Exception ex) { logger.Error("处理流程异常终止", ex); throw; } }
内容的提问来源于stack exchange,提问作者Mike
相关产品推荐
相关产品推荐

