【问题标题】:Fix for C# Async Method Executing Sequentially修复了按顺序执行的 C# 异步方法
【发布时间】:2019-04-04 20:55:47
【问题描述】:

我遇到的问题是对我正在执行的异步方法的调用是按顺序发生的。我正在将调用的任务添加到 ConcurrentBag 并等待包中的任务。我不关心这些调用的结果,我只需要确认它们已完成。但是,这些调用是完全按顺序发生的,这非常令人困惑。有问题的方法通过带有参数化查询的 Npgsql 执行一些 PostgreSQL 查询。调用者获取一棵我们自己的数据树,并拉出树中的所有节点并遍历节点并对它们执行此任务。我还使用了一个自定义 AsyncHelper 类,它将遍历 IEnumerable 实现器中的任务并等待其中的任务。我的 Tree 实现和 AsyncHelper 都在另一段代码中进行了测试,该代码执行与此代码相同的基本原则,它按预期异步执行任务。

我添加了对函数调用的日志记录,以确认这些是按顺序发生的。我还从包中取出方法并运行该方法,它仍然做同样的事情,它按顺序发生并且在完成之前不会继续我的循环。我所有的方法都被标记为异步,直到循环之后我才等待它们。

//method executing sequentially
public static async Task<List<ContactStatistic>> getContactStats(Guid tenantId, DateTime start, DateTime end, Breakdown breakdown) {
    if (!await Postgres.warmConnection(5)) { return null; }
    var hierarchy = await getTreeForTenant<TenantContactStatsNode>(tenantId);

    //perform calculations to determine stats for each element
    var calculationTasks = new ConcurrentBag<Task>();
    var allData = await hierarchy.getAllData();
    var timestampGotAllData = DateTime.Now;

    foreach (var d in allData) {
        calculationTasks.Add(d.getContactStats(start, end, breakdown));
    }

    Console.WriteLine("about to await all the tasks");
    //await the tasks to complete for calculations
    await AsyncHelper.waitAll(calculationTasks);
}


//method it's calling
public async Task getContactStats(DateTime start, DateTime end, Breakdown breakdown) {
    //perform two async postgres calls
    //await postgres calls
    //validate PG response
    //perform manipluation on this object with data from the queries
}

我希望第一次调用调用第二个函数,将任务添加到包中,并在完成后等待它们。实际发生的是该方法正在运行、完成,然后添加到包中。

* 编辑 *

以下是第二次调用的完整代码。它根据时间从数据库中获取一些数据,填补回退时间之间的空白,因此我们有一个完全顺序的返回列表,包括数据库中没有数据的所有时间,并将其放入对象级变量中

public async Task getContactStats(DateTime start, DateTime end, Breakdown breakdown) {
    if (breakdown == Breakdown.Month) {
        //max out month start day to include all for the initial month in the initial count
        start = new DateTime(start.Year, start.Month, DateTime.DaysInMonth(start.Year, start.Month));
    } else {
        //day breakdown previous stats should start the day before given start day
        start = start.AddDays(-1);
    }

    var tran = new PgTran();
    var breakdownQuery = breakdown == Breakdown.Day ? Queries.GET_CONTACT_DAY_BREAKDOWN : Queries.GET_CONTACT_MONTH_BREAKDOWN;
    tran.setQueries(Queries.GET_CONTACT_COUNT_BEFORE_DATE, breakdownQuery);
    tran.setParams(new NpgsqlParameter("@tid", tenantId), new NpgsqlParameter("@start", start), new NpgsqlParameter("@end", end));
    var tranResults = await Postgres.getAll<ContactDayStatistic>(tran);
    //ensure transaction returns two query results
    if (tranResults == null || tranResults.Count != 2) { return; }


    //ensure valid past count was retrieved
    var prevCountResult = tranResults[0];
    if (prevCountResult == null || prevCountResult.Count != 1) { return; }
    var prevStat = new ContactDayStatistic(start.Day, start.Month, start.Year, prevCountResult[0].count);
    //ensure valid contact stat breakdown was retrieved
    var statBreakdown = tranResults[1];
    if (statBreakdown == null) { return;}

    var datesInBreakdown = new List<DateTime?>();
    //get all dates in the returned stats
    foreach (var e in statBreakdown) {
        var eventDate = new DateTime(e.year, e.month, e.day);
        if (datesInBreakdown.Find(item => item == eventDate) == null)
            datesInBreakdown.Add(eventDate);
    }
    //sort so they are sequential
    datesInBreakdown.Sort();

    //initialize timeline starting with initial breakdown
    var fullTimeline = new List<ContactStatistic>();
    //convert initial stat to the right type for final display
    fullTimeline.Add(breakdown == Breakdown.Month ? new ContactStatistic(prevStat) : prevStat);
    foreach (var d in datesInBreakdown) {
        //null date is useless, won't occur, nullable date just for default value of null
        if (d == null) { continue; }
        var newDate = d.Value;
        //fill gaps between last date given and this date
        ContactStatistic.fillGaps(breakdown, newDate, prevStat.getDate(), prevStat.count, ref fullTimeline, false);
        //get stat for this day
        var stat = statBreakdown.Find(item => d == new DateTime(item.year, item.month, item.day));
        if (stat == null) { continue; }
        //add last total for a rolling total of count
        stat.count += prevStat.count;
        fullTimeline.Add(breakdown == Breakdown.Month ? new ContactStatistic(stat) : stat);
        prevStat = stat;
    }
    //fill gaps between last date and end
    ContactStatistic.fillGaps(breakdown, end, prevStat.getDate(), prevStat.count, ref fullTimeline, true);
    //cast list to appropriate return type
    contactStats.Clear();
    contactStats = fullTimeline;
}

* 编辑 2 * 这是 AsyncHelper 用于等待这些任务的代码。这个函数非常适合使用相同框架的其他代码,它基本上只是清理必须等待枚举任务的代码。

public static async Task waitAll(IEnumerable<Task> coll) {
    foreach (var taskToWait in coll) {
        await taskToWait;
    }  
}

* 编辑 3 * 根据建议,我将 waitAll() 更改为使用 Task.WhenAll() 而不是 foreach 循环,但是问题仍然存在。

public static async Task waitAll(IEnumerable<Task> coll) {
    await Task.WhenAll(coll);
}

* 编辑 4 * 为了确保不是 Postgres 调用导致这种情况发生,我将第二种方法更改为仅执行打印行,然后休眠 200 毫秒以保持执行路径清晰。我仍然注意到这是完全按顺序发生的(甚至导致我对该函数的 POST 超时,因为实际的实际调用需要将近 20 毫秒)。下面是演示该更改的代码

public async Task getContactStats(DateTime start, DateTime end, Breakdown breakdown) {
    Console.WriteLine("CALLED!");
    Thread.Sleep(200);
}

* 编辑 5 * 根据建议,我尝试了一个并行的 foreach 来尝试填充任务的 ConcurrentBag 而不是普通的 foreach。我在这里遇到了一个问题,即并行 foreach 会在第一次添加完成后完成,并且不会一次添加所有任务。

var calculationTasks = new ConcurrentBag<Task>();
var allData = await hierarchy.getAllData();
var timestampGotAllData = DateTime.Now;
Parallel.ForEach(allData, item => {
    Console.WriteLine("trying parallel foreach");
    calculationTasks.Add(item.getContactStats(start, end, breakdown));
});

Console.WriteLine("about to await all the tasks");
//await the tasks to complete for calculations
await AsyncHelper.waitAll(calculationTasks);

* 编辑 6 * 为了视觉,我运行了代码并做了一些输出来显示正在发生的怪事。执行代码如下:

foreach (var d in allData) {
    Console.WriteLine("Adding call to bag");
    calculationTasks.Add(d.getContactStats(start, end, breakdown));
    Console.WriteLine("Done adding call to bag");
}

输出为:https://i.imgur.com/3y5S4eS.png

因为它每次都打印“CALLED”,所以“Done!”在“完成对包的调用”之前,这些执行是按顺序发生的,而不是像预期的那样异步。

【问题讨论】:

  • 发布第二次调用的代码。
  • @EvanTrimboli 编辑以包含第二次调用的完整代码
  • 您写道,您的异步助手正在等待里面的所有任务。那个帮手究竟是怎么等他们的?是否明确执行 .Wait 或类似操作?我怀疑这是你问题的根源
  • @AlfredoMS 我为 AsyncHelper 的方法添加了代码来执行此操作,它基本上是语法糖,可以更好地等待枚举任务。我已经在没有 ConcurrentBag 的情况下测试了同样的调用,只执行了该方法,但它仍然是连续的
  • pgTran 中的代码是什么?

标签: c# postgresql lambda async-await


【解决方案1】:

我对此的直觉是,这将与您在方法中打开的交易有关。由于这里似乎有一些自定义类,因此很难准确判断代码中发生了什么 - 但是当您打开事务时是否可能会发生一些锁定?由于这发生在您第一次等待之前,它必须在等待代码之前“按顺序”运行。

您的自定义 'waitall' 方法似乎不是问题,但您应该考虑删除它并使用内置的 Task.WhenAll 来异步等待这些。

【讨论】:

  • 事务本身没有任何等待,对 Postgres 的调用按预期发生并且正确异步,我在其他地方使用这些确切的调用并且它们按预期工作。它似乎植根于第二种方法的调用,就像第二种方法只是拒绝异步执行一样。对于您的第二个建议,我根据您和 Mike 的建议将 AsyncHelper 更改为使用 Task.WhenAll() 而不是 foreach 循环,但问题仍然存在
  • 为了对我正在使用的 Postgres 类进行一些额外的说明,事务在一次访问 Postgres 的过程中构建并执行了两个查询。我正在使用 Npgsql 进行连接,它使用池来允许并发操作。但是,即使我完全删除 Postgres 调用并在调用中执行 Console.WriteLine() 和 Sleep,它仍然会同步发生。
  • 你用什么来睡觉,你把它放在代码的什么地方? (如果是在第一次调用await之前,那就是同步休眠。你应该在等待Task.Delay(xxx)。Console.WriteLine也会排队,只能同步写入:docs.microsoft.com/en-us/dotnet/api/…
  • 好点,我确实将 Sleep 切换为 Task.Delay 并再次运行,仍然是连续的。
  • @ErikA - 您是否等待 task.delay 并将 console.writeline 移到延迟后?
【解决方案2】:

试试这个:

foreach (var d in allData) 
{
    calculationTasks.Add(Task.Run(() => d.getContactStats(start, end, breakdown)));
}

//Other code here
//...

Task.WaitAll(calculationTasks.ToArray());

我们实际上是在创建一个“运行”您的方法的任务。然后我们等待这些任务完成。

诚然,我不完全确定为什么您的版本会阻塞,但这似乎可以解决问题。

更新:

我通过输出线程 id 进行测试,并且 OP 的版本在同一个线程上执行任务。也许线程被包锁住了,这迫使新任务等待?我提出的解决方案在不同的线程 ID 上产生了结果,我认为这可以解释为什么它不会阻塞。

【讨论】:

  • 好主意!,它应该工作是有道理的,因为任务本身不应该等待,它是异步的并且不等待所以它真的不应该阻塞,但遗憾的是它仍然在尝试阻塞: '(
  • @ErikA:我在本地对此进行了测试,它对我有用。你改成WaitAll而不是你的助手了吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-02-24
  • 1970-01-01
  • 2021-09-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多