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

为何我的ActionBlock未手动设置就进入Completed状态?

问题描述

运行约15分钟后,ParallelWorker.SendAsync(rd)返回false,且ParallelWorker.Complete为true,触发异常。已确认ProcessSetupsAsync内部同步代码未抛出异常(catch块从未执行)。

场景:通过静态ActionBlock处理大量生成的工作项,更换数据集可复现该问题。

原始代码

public static class Setups
{
    private struct RunData
    {
        internal MyClass1 Setup;
        internal MyClass2 Positions;
        internal string Set;
    }

    private static readonly ActionBlock<RunData> ParallelWorker = 
        new(d => ProcessSetupsAsync(d.Setup, d.Positions, d.Set), 
            new ExecutionDataflowBlockOptions 
            { 
                BoundedCapacity = Environment.ProcessorCount * 10, 
                MaxDegreeOfParallelism = Environment.ProcessorCount, 
                SingleProducerConstrained = false 
            });

    public static async Task GetSetups(ItemType[] bar1Filter, 
        ItemType[] bar2Filter, bool sameSet)
    {
        for (/* do some work*/)
        {
            foreach (RunData rd in MyMethod1(/*variables*/).Select(final 
                => new RunData { Positions = positions, Set = set, Setup = final }))
            {
                if (!await ParallelWorker.SendAsync(rd).ConfigureAwait(false))
                {
                    // Fails after running for about 15 minutes
                    // with ParallelWorker.Complete is true
                    throw new Exception("xxxxxxxx");
                }
            }
        }
        ParallelWorker.Complete();
        await ParallelWorker.Completion.ConfigureAwait(false);
    }

    private static Task ProcessSetupsAsync
                        (MyClass1 setup, MyClass2 positions, string set)
    {
        try
        {
            List<MyClass1> setups = new();

            /* Do some work */
            return setups.Count > 0 
                ? LocalAccess.AddSetupAsync(setups) 
                : Task.CompletedTask;
        }
        catch (Exception ex)
        {
            Console.WriteLine(ex); //<-Never gets hit.
            throw;
        }
    }
}
问题分析与解决方案

核心原因

  1. 异步任务异常未被捕获:ProcessSetupsAsync的try/catch仅覆盖同步代码逻辑,LocalAccess.AddSetupAsync返回的异步任务若抛出异常,属于异步异常,不会被当前捕获逻辑处理。此时ActionBlock会因委托任务故障进入完成状态,拒绝接收新消息。
  2. 静态实例被意外终止:ParallelWorker是静态字段,若应用内其他代码路径调用了ParallelWorker.Complete(),会导致块提前进入完成状态。
  3. 故障状态未被感知:ActionBlock进入故障状态后,Completion任务会包含异常,但原代码未捕获这部分异常,无法定位具体错误根源。

解决步骤

  • 修复异步异常捕获逻辑:将ProcessSetupsAsync改为async Task,通过await捕获异步任务的异常:
    private static async Task ProcessSetupsAsync(MyClass1 setup, MyClass2 positions, string set)
    {
        try
        {
            List<MyClass1> setups = new();
    
            /* Do some work */
            if (setups.Count > 0)
            {
                await LocalAccess.AddSetupAsync(setups).ConfigureAwait(false);
            }
        }
        catch (Exception ex)
        {
            Console.WriteLine(ex);
            throw;
        }
    }
    
  • 避免静态实例的状态污染:将ParallelWorker改为非静态,每次调用GetSetups时创建新实例,防止跨调用的状态干扰:
    public static async Task GetSetups(ItemType[] bar1Filter, 
        ItemType[] bar2Filter, bool sameSet)
    {
        var parallelWorker = new ActionBlock<RunData>(d => ProcessSetupsAsync(d.Setup, d.Positions, d.Set), 
            new ExecutionDataflowBlockOptions 
            { 
                BoundedCapacity = Environment.ProcessorCount * 10, 
                MaxDegreeOfParallelism = Environment.ProcessorCount, 
                SingleProducerConstrained = false 
            });
    
        for (/* do some work*/)
        {
            foreach (RunData rd in MyMethod1(/*variables*/).Select(final 
                => new RunData { Positions = positions, Set = set, Setup = final }))
            {
                if (!await parallelWorker.SendAsync(rd).ConfigureAwait(false))
                {
                    throw new Exception("xxxxxxxx");
                }
            }
        }
        parallelWorker.Complete();
        await parallelWorker.Completion.ConfigureAwait(false);
    }
    
  • 捕获ActionBlock的故障异常:在等待Completion时捕获AggregateException,定位具体错误:
    try
    {
        await parallelWorker.Completion.ConfigureAwait(false);
    }
    catch (AggregateException ex)
    {
        foreach (var innerEx in ex.InnerExceptions)
        {
            Console.WriteLine($"ActionBlock故障原因:{innerEx.Message}\n{innerEx.StackTrace}");
        }
        throw;
    }
    

内容的提问来源于stack exchange,提问作者Iaman Swtrse

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:37:23