【问题标题】:Can you avoid Task Continuations when using async/await if you want execution to continue immediately如果您想立即继续执行,您可以在使用 async/await 时避免任务继续吗
【发布时间】:2020-12-03 19:17:16
【问题描述】:

我正在研究一个协议,并尝试尽可能多地使用 async/await 以使其能够很好地扩展。该协议必须支持成百上千的同时连接。下面是一点点伪代码来说明我的问题。

private static async void DoSomeWork()
{
    var protocol = new FooProtocol();
    await protocol.Connect("127.0.0.1", 1234);
    var i = 0;
    while(i != int.MaxValue)
    {
        i++;
        var request = new FooRequest();
        request.Payload = "Request Nr " + i;
        var task = protocol.Send(request);
        _ = task.ContinueWith(async tmp =>
        {
            var resp = await task;
            Console.WriteLine($"Request {resp.SequenceNr} Successful: {(resp.Status == 0)}");
         });

    }
}

下面是协议的一些伪代码。

public class FooProtocol
{
    private int sequenceNr = 0;
    
    private SemaphoreSlim ss = new SemaphoreSlim(20, 20);
    
    public Task<FooResponse> Send(FooRequest fooRequest)
    {
        var tcs = new TaskCompletionSource<FooResponse>();
        ss.Wait();
        var tmp = Interlocked.Increment(ref sequenceNr);
        fooRequest.SequenceNr = tmp;
        // Faking some arbitrary delay. This work is done over sockets. 
        Task.Run(async () =>
        {
            await Task.Delay(1000);
            tcs.SetResult(new FooResponse() {SequenceNr = tmp});
            ss.Release();
        });
        return tcs.Task;

    }
}

我有一个包含请求和响应对的协议。我使用过异步套接字编程。 FooProtocol 将负责将请求与响应(序列号)进行匹配,并且还将负责挂起请求的最大数量。 (在伪和我的代码中使用信号量苗条完成,所以我不担心请求失控)。 DoSomeWork 方法调用 Protocol.Send 方法,但我不想等待响应,我想旋转并发送下一个,直到我被最大数量的未决请求阻塞。当任务完成后,我想检查响应并可能做一些工作。

我想解决两件事

  1. 我想避免使用 Task.ContinueWith(),因为它似乎不完全符合 async/await 模式
  2. 因为我一直在等待连接,所以我不得不使用 async 修饰符。现在我收到来自 IDE 的警告“因为不等待此调用,所以在此调用完成之前继续执行当前方法。考虑将 'await' 运算符应用于调用结果。”我不想这样做,因为一旦我这样做,它就会破坏协议处理许多请求的能力。我可以摆脱警告的唯一方法是使用丢弃。这不是最糟糕的事情,但我不禁觉得我错过了一个技巧,并且太努力了。

【问题讨论】:

  • 您可以将作业委托给线程池并忘记它 -> _ = Task.Run(async () =&gt; { var resp = await protocol.Send(request); Console.WriteLine($"Request {resp.SequenceNr} Successful: {(resp.Status == 0)}"); }
  • @CamiloTerevinto 听起来有成百上千的并发请求,你会饿死线程池。除此之外,它有点违背异步的目的。
  • @GSerg 是的,我担心会是这样。也许只是推送发送请求并获得对async void 方法的响应?
  • 将await protocol.Send(request); Console.WriteLine 分解为单独的Func&lt;Task&gt;,继续创建它们并将它们添加到列表中,然后最后对它们执行Task.WhenAll?再说一次,你可能想要实现某种throttling。
  • @GSerg Throttling 内置于协议中。我把它省略了伪代码。我无法保留所有任务的列表,因为这些连接是长期存在的,可能在服务的生命周期内。

标签: c# asynchronous concurrency


【解决方案1】:

旁注:我希望您的实际代码使用SemaphoreSlim.WaitAsync 而不是SemaphoreSlim.Wait。

在大多数套接字代码中,您最终会得到一个连接列表,并且每个连接都有某种“处理器”。在异步世界中,这自然表示为 Task。

因此您需要保留Tasks 的列表;至少,您的消费应用程序需要知道何时可以安全关闭(即,已收到所有响应)。

不要急于使用Task.Run;只要你没有阻塞(例如,SemaphoreSlim.Wait),你可能不会饿死线程池。请记住,在awaits 期间,没有使用线程池线程。

【讨论】:

  • 是的。我正在从内存中写下代码。我确实有一个连接字典,以便我可以查找请求的目标连接。 WRT 到任务列表,这是一个长期运行的服务并保留一个列表,它将变得庞大。我确实想到了某种数据结构,我可以定期清理已完成的任务。你会推荐什么?
  • 我确实注意到我的 SemaphoreSlim WaitAsyc 尝试执行 await ss.WaitAsync() 然后说我必须让 Send 方法返回 Task>。我可以摆脱它的唯一方法是将它包装在另一个类或方法中,然后使用一个非常禁忌的方法 public async void。在协议类本身中,我使用的是 TaskCompletionSource 字典。交回任务,然后以另一种方法异步侦听套接字,我查找 TCS 并设置结果。有没有其他方法可以做到这一点?
  • @uriDium:一旦连接不再可行,任务应通知“连接管理器”以删除它们,因此您不应该有任何非活动集合的任务。关于异步Send,你应该可以用async/await/Task来做到这一点。您可以从 TCS 控制 Task,您可以在 WaitAsync 之后打开 await。
【解决方案2】:

我不确定在协议级别强制执行最大并发性是否是个好主意。在我看来,这个责任属于协议的调用者。所以我会删除SemaphoreSlim,让它做一件它知道做得好的事情:

public class FooProtocol
{
    private int sequenceNr = 0;

    public async Task<FooResponse> Send(FooRequest fooRequest)
    {
        var tmp = Interlocked.Increment(ref sequenceNr);
        fooRequest.SequenceNr = tmp;
        await Task.Delay(1000); // Faking some arbitrary delay
        return new FooResponse() { SequenceNr = tmp };
    }
}

然后我将使用来自TPL Dataflow 库的ActionBlock 来协调通过协议发送大量请求的过程,通过处理并发、背压(BoundedCapacity)、取消(如果需要)、错误处理和整个操作的状态(运行、完成、失败等)。示例:

private static async Task DoSomeWorkAsync()
{
    var protocol = new FooProtocol();

    var actionBlock = new ActionBlock<FooRequest>(async request =>
    {
        var resp = await protocol.Send(request);
        Console.WriteLine($"Request {resp.SequenceNr} Status: {resp.Status}");
    }, new ExecutionDataflowBlockOptions()
    {
        MaxDegreeOfParallelism = 20,
        BoundedCapacity = 100
    });

    await protocol.Connect("127.0.0.1", 1234);

    foreach (var i in Enumerable.Range(0, Int32.MaxValue))
    {
        var request = new FooRequest();
        request.Payload = "Request Nr " + i;
        var accepted = await actionBlock.SendAsync(request);
        if (!accepted) break; // The block has failed irrecoverably
    }
    actionBlock.Complete();
    await actionBlock.Completion; // Propagate any exceptions
}

BoundedCapacity = 100 配置意味着ActionBlock 将在其内部缓冲区中最多存储 100 个请求。当达到这个阈值时,任何想要向它发送更多请求的人都必须等待。等待将发生在await actionBlock.SendAsync 行中。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-05-24
    • 1970-01-01
    • 1970-01-01
    • 2013-05-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-04-08
    相关资源
    最近更新 更多