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

并行生产者-消费者模式下,如何让消费者感知生产者异常?

解决生产者异常时消费者感知问题的几种可行方案

方案一:追踪生产者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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 23:10:22