【问题标题】:C# async within an action动作中的 C# 异步
【发布时间】:2016-08-28 13:31:23
【问题描述】:

我想写一个接受多个参数的方法,包括一个动作和一个重试量并调用它。

所以我有这个代码:

public static IEnumerable<Task> RunWithRetries<T>(List<T> source, int threads, Func<T, Task<bool>> action, int retries, string method)
    {
        object lockObj = new object();
        int index = 0;

        return new Action(async () =>
        {
            while (true)
            {
                T item;
                lock (lockObj)
                {
                    if (index < source.Count)
                    {
                        item = source[index];
                        index++;
                    }
                    else
                        break;
                }

                int retry = retries;
                while (retry > 0)
                {
                    try
                    {
                        bool res = await action(item);
                        if (res)
                            retry = -1;
                        else
                            //sleep if not success..
                            Thread.Sleep(200);

                    }
                    catch (Exception e)
                    {
                        LoggerAgent.LogException(e, method);
                    }
                    finally
                    {
                        retry--;
                    }
                }
            }
        }).RunParallel(threads);
    }

RunParallel 是 Action 的扩展方法,如下所示:

public static IEnumerable<Task> RunParallel(this Action action, int amount)
    {
        List<Task> tasks = new List<Task>();
        for (int i = 0; i < amount; i++)
        {
            Task task = Task.Factory.StartNew(action);
            tasks.Add(task);
        }
        return tasks;
    }

现在,问题是:线程只是在没有等待操作完成的情况下消失或崩溃。

我写了这个示例代码:

private static async Task ex()
    {
        List<int> ints = new List<int>();
        for (int i = 0; i < 1000; i++)
        {
            ints.Add(i);
        }

        var tasks = RetryComponent.RunWithRetries(ints, 100, async (num) =>
        {
            try
            {
                List<string> test = await fetchSmthFromDb();
                Console.WriteLine("#" + num + "  " + test[0]);
                return test[0] == "test";
            }
            catch (Exception e)
            {
                Console.WriteLine(e.StackTrace);
                return false;
            }

        }, 5, "test");

        await Task.WhenAll(tasks);
    }

fetchSmthFromDb 是一个简单的任务>,它从数据库中获取一些东西,并且在本示例之外调用时可以正常工作。

每当调用List&lt;string&gt; test = await fetchSmthFromDb(); 行时,线程似乎正在关闭并且Console.WriteLine("#" + num + " " + test[0]); 甚至没有被触发,在调试断点时也从未命中。

最终工作代码

private static async Task DoWithRetries(Func<Task> action, int retryCount, string method)
    {
        while (true)
        {
            try
            {
                await action();
                break;
            }
            catch (Exception e)
            {
                LoggerAgent.LogException(e, method);
            }

            if (retryCount <= 0)
                break;

            retryCount--;
            await Task.Delay(200);
        };
    }

    public static async Task RunWithRetries<T>(List<T> source, int threads, Func<T, Task<bool>> action, int retries, string method)
    {
        Func<T, Task> newAction = async (item) =>
        {
            await DoWithRetries(async ()=>
            {
                await action(item);
            }, retries, method);
        };
        await source.ParallelForEachAsync(newAction, threads);
    }

【问题讨论】:

  • 您确定您的记录器是线程安全的吗?当我用 Console.WriteLine 替换它时,我得到“线程被中止”......还有什么是锁?你想做什么?
  • 我真的对上面的例子感到困惑。为什么要尝试并行运行相同的操作 100 次? (RunParallel 方法)是对数据库进行某种负载测试吗?
  • Logger 不是线程安全的,但对我来说不会崩溃。 @SergeSemenov 我正在使用 mongodb,我无法像在 SQL 中那样在一个过程中更新 100 个文件,所以我构建了一个方法来接受单个可枚举的操作列表并作为单个过程操作
  • 你做错了,因为你运行你的 while (true) 循环 100 次
  • 我很乐意提供见解

标签: c# task action


【解决方案1】:

问题出在这一行:

return new Action(async () => ...

您使用 async lambda 启动异步操作,但不返回要等待的任务。 IE。它在工作线程上运行,但你永远不会知道它何时完成。并且您的程序在异步操作完成之前终止 - 这就是您看不到任何输出的原因。

必须是:

return new Func<Task>(async () => ...

更新

首先,您需要拆分方法的职责,因此不要将重试策略(不应硬编码为检查布尔结果)与并行运行的任务混合。

然后,如前所述,您运行 while (true) 循环 100 次,而不是并行执行。

正如@MachineLearning 指出的那样,使用Task.Delay 而不是Thread.Sleep

总体而言,您的解决方案如下所示:

using System.Collections.Async;

static async Task DoWithRetries(Func<Task> action, int retryCount, string method)
{
    while (true)
    {
        try
        {
            await action();
            break;
        }
        catch (Exception e)
        {
            LoggerAgent.LogException(e, method);
        }

        if (retryCount <= 0)
            break;

        retryCount--;
        await Task.Delay(millisecondsDelay: 200);
    };
}

static async Task Example()
{
    List<int> ints = new List<int>();
    for (int i = 0; i < 1000; i++)
        ints.Add(i);

    Func<int, Task> actionOnItem =
        async item =>
        {
            await DoWithRetries(async () =>
            {
                List<string> test = await fetchSmthFromDb();
                Console.WriteLine("#" + item + "  " + test[0]);
                if (test[0] != "test")
                    throw new InvalidOperationException("unexpected result"); // will be re-tried
            },
            retryCount: 5,
            method: "test");
        };

    await ints.ParallelForEachAsync(actionOnItem, maxDegreeOfParalellism: 100);
}

您需要使用AsyncEnumerator NuGet Package 才能使用System.Collections.Async 命名空间中的ParallelForEachAsync 扩展方法。

【讨论】:

  • 感谢重播,但是我该如何处理扩展方法,为 Fun 创建另一个扩展?
  • 我建议你用你改变的代码更新你的问题
  • 但是为什么线程内的 DoWithRetries 呢?如果我在每个方法中编写重试机制,我已经可以解决它。我想要一些东西来包装原始任务而不是破坏它。我需要将线程数量用作输入,目的是保存所有这些包含主任务的代码,而不是复制粘贴它
  • 感谢您在回答中提及我的建议! :-)
  • 两者都没有,我修改了代码,使其更通用且更有效。非常感谢
【解决方案2】:

除了最终的完全重新设计之外,我认为强调原始代码的真正错误非常重要。

0) 首先,正如@Serge Semenov 立即指出的那样,必须将 Action 替换为

Func<Task>

但还有其他两个重要的变化。

1) 使用异步委托作为参数,必须使用更新的 Task.Run 而不是旧的模式 new TaskFactory.StartNew(否则您必须显式添加 Unwrap())

2) 此外,ex() 方法不能是异步的,因为 Task.WhenAll 必须使用 Wait() 等待,而无需等待。

到那时,即使存在需要重新设计的逻辑错误,但从纯技术角度来看,它确实有效并且产生了输出。

在线测试:http://rextester.com/HMMI93124

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-16
    • 2019-01-11
    • 2012-04-14
    • 1970-01-01
    • 1970-01-01
    • 2011-05-02
    相关资源
    最近更新 更多