【问题标题】:Mixing async tasks with blocking sync task混合异步任务和阻塞同步任务
【发布时间】:2015-07-09 23:09:20
【问题描述】:

我正在编写一组异步任务来消除下载和解析数据,但是我在下一步更新数据库时遇到了一些空白。

问题是,为了提高性能,我使用 TableLock 来加载相当大的数据集,所以我想做的是让我的导入服务等待第一个任务返回,然后开始导入。如果另一个任务在第一次导入运行时完成,该进程将加入队列并等待任务 1 的导入服务完成。

例如。

异步 - 任务1 - 任务2 - 任务3

同步 - 进口服务

RunAsync 任务

Task3 returns first > ImportService.Import(Task3)
Task1 return, ImportService is still running. Wait()
ImportService.Complete() event
Task2 returns. Wait()
ImportService.Import(Task1)
ImportService.Complete() event
ImportService.Import(Task2)
ImportService.Complete() event

希望这是有道理的!

【问题讨论】:

  • 你可能应该考虑TPL DataFlow
  • 保罗,这正是我想要的!谢谢!

标签: c# asynchronous async-await


【解决方案1】:

你不能在这里真正使用 await,但你可以等待多个任务完成:

var tasks = new List<Task)();
// start the tasks however
tasks.Add(Task.Run(Task1Function);
tasks.Add(Task.Run(Task2Function);
tasks.Add(Task.Run(Task2Function);

while (tasks.Count > 0)
{
   var i = Task.WaitAny(tasks.ToArray()); // yes this is ugly but an array is required
   var task = tasks[i];
   tasks.RemoveAt(i);
   ImportService.Import(task); // do you need to pass the task or the task.Result
}

在我看来,应该有更好的选择。例如,您可以让任务和导入运行并在 ImportService 部分添加锁:

// This is the task code doing whatever
....
// Task finishes and calls ImportService.Import
lock(typeof(ImportService)) // actually the lock should probably be inside the Import method
{
   ImportService.Import(....);
}

您的要求有几件事困扰着我(包括使用静态 ImportService,静态类很少是个好主意),但如果没有更多详细信息,我无法提供更好的建议。

【讨论】:

  • 致 OP:我同意 Eli 的观点,即您应该简单地同步导入,每个任务在导入之前获取锁作为他们执行的最后一个操作。根据具体的实现,这个主题可能会有更好的变化,但上面的例子是可以合理提供的,除非你更详细地改进问题,包括清楚地说明问题的a good, minimal, complete code example。跨度>
  • 谢谢伊莱和彼得。我非常喜欢锁定线程的想法,这是我最初的想法,但不认为使用 Tasks 是可能的,你每天都会学到新的东西。谢谢! :)
  • @Al.您应该考虑是否需要在应用程序级别进行锁定。我不知道您使用的是什么数据库,但数据库擅长的一件事是处理并发操作。如果您在应用程序级别锁定,假设甚至需要锁定,那么您将被限制为同时运行的单个应用程序。无需重写锁定机制就无法扩展到多台服务器。也许您的问题应该是:如何在没有 TableLock 或没有应用程序级锁的情况下运行我的查询。
【解决方案2】:

虽然这可能不是最优雅的解决方案,但我会尝试启动工作任务并将它们的输出放在 ConcurrentQueue 中。您可以在计时器上检查队列中的工作,直到所有任务都完成。

var rand = new Random();
var importedData = new List<string>();
var results = new ConcurrentQueue<string>();
var tasks = new List<Task<string>>
{
    new Task<string>(() =>
    {
        Thread.Sleep(rand.Next(1000, 5000));
        Debug.WriteLine("Task 1 Completed");
        return "ABC";
    }),
    new Task<string>(() =>
    {
        Thread.Sleep(rand.Next(1000, 5000));
        Debug.WriteLine("Task 2 Completed");
        return "FOO";
    }),
    new Task<string>(() =>
    {
        Thread.Sleep(rand.Next(1000, 5000));
        Debug.WriteLine("Task 3 Completed");
        return "BAR";
    })
};

tasks.ForEach(t =>
{
    t.ContinueWith(r => results.Enqueue(r.Result));
    t.Start();
});

var allTasksCompleted = new AutoResetEvent(false);
new Timer(state =>
{
    var timer = (Timer) state;
    string item;

    if (!results.TryDequeue(out item)) 
        return;

    importedData.Add(item);
    Debug.WriteLine("Imported " + item);

    if (importedData.Count == tasks.Count)
    {
        timer.Dispose();
        Debug.WriteLine("Completed.");
        allTasksCompleted.Set();
    }
}).Change(1000, 100);


allTasksCompleted.WaitOne();

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-07-07
    • 1970-01-01
    • 2011-05-07
    • 1970-01-01
    • 2019-01-24
    • 2012-09-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多