【问题标题】:Network Command Processing with TPL Dataflow使用 TPL 数据流处理网络命令
【发布时间】:2014-01-31 16:32:43
【问题描述】:

我正在开发一个系统,该系统涉及通过 TCP 网络连接接受命令,然后在执行这些命令时发送响应。相当基本的东西,但我希望支持一些要求:

  1. 多个客户端可以同时连接并建立单独的会话。会话可以根据需要持续多久,如果需要,相同的客户端 IP 可以建立多个并行会话。
  2. 每个会话可以同时处理多个命令,因为某些请求的操作可以并行执行。

我想使用 async/await 干净利落地实现这一点,根据我所读到的内容,TPL Dataflow 听起来是一种很好的方法,可以将处理干净地分解成可以在线程池上运行的好块,而不是在线程池上运行为不同的会话/命令绑定线程,阻塞等待句柄。

这就是我要开始的内容(为了简化,去掉了一些部分,例如异常处理的细节;我还省略了一个为网络 I/O 提供高效等待的包装器):

    private readonly Task _serviceTask;
    private readonly Task _commandsTask;
    private readonly CancellationTokenSource _cancellation;
    private readonly BufferBlock<Command> _pendingCommands;

    public NetworkService(ICommandProcessor commandProcessor)
    {
        _commandProcessor = commandProcessor;
        IsRunning = true;
        _cancellation = new CancellationTokenSource();
        _pendingCommands = new BufferBlock<Command>();
        _serviceTask = Task.Run((Func<Task>)RunService);
        _commandsTask = Task.Run((Func<Task>)RunCommands);
    }

    public bool IsRunning { get; private set; }

    private async Task RunService()
    {
        _listener = new TcpListener(IPAddress.Any, ServicePort);
        _listener.Start();

        while (IsRunning)
        {
            Socket client = null;
            try
            {
                client = await _listener.AcceptSocketAsync();
                client.Blocking = false;

                var session = RunSession(client);
                lock (_sessions)
                {
                    _sessions.Add(session);
                }
            }
            catch (Exception ex)
            { //Handling here...
            }
        }
    }

    private async Task RunCommands()
    {
        while (IsRunning)
        {
            var command = await _pendingCommands.ReceiveAsync(_cancellation.Token);
            var task = Task.Run(() => RunCommand(command));
        }
    }

    private async Task RunCommand(Command command)
    {
        try
        {
            var response = await _commandProcessor.RunCommand(command.Content);
            Send(command.Client, response);
        }
        catch (Exception ex)
        {
            //Deal with general command exceptions here...
        }
    }

    private async Task RunSession(Socket client)
    {
        while (client.Connected)
        {
            var reader = new DelimitedCommandReader(client);

            try
            {
                var content = await reader.ReceiveCommand();
                _pendingCommands.Post(new Command(client, content));
            }
            catch (Exception ex)
            {
                //Exception handling here...
            }
        }
    }

基础知识看起来很​​简单,但有一部分让我感到困惑:如何确保在关闭应用程序时等待所有待处理的命令任务完成?当我使用 Task.Run 执行命令时,我得到了 Task 对象,但是我如何跟踪待处理的命令,以便在允许服务关闭之前确保所有这些命令都已完成?

我考虑过使用简单的列表,并在完成时从列表中删除命令,但我想知道我是否缺少 TPL 数据流中的一些基本工具,这些工具可以让我更干净地完成这项工作。


编辑:

阅读有关 TPL 数据流的更多信息,我想知道我是否应该使用 TransformBlock 并增加 MaxDegreeOfParallelism 以允许处理并行命令?这对可以并行运行的命令数量设置了上限,但我认为这对我的系统来说是一个合理的限制。我很想听听那些有 TPL Dataflow 经验的人来了解我是否走在正确的轨道上。

【问题讨论】:

  • 我强烈建议您使用自托管 WebAPI 而不是 TCP/IP 服务器,除非您正在与没有 HTTP 客户端组件的嵌入式机器通信。裸 TCP/IP 有 的陷阱... Socket.Connected 不能像您在此处尝试使用的那样使用。即使您从未收到命令,您也需要发送周期性数据以避免半开问题。关机总是很难做到正确。简而言之,编写一个正确的 TCP/IP 服务器只是真的很难
  • @StephenCleary,我正在与一个仅本机支持原始 TCP 协议的机器人集成。我同意这方面存在挑战,并且此 sn-p 不包括针对发生的各种断开连接情况的所有异常处理,以及会话超时处理。
  • 你用Socket.Blocking = false做什么?很长一段时间我都没有看到非阻塞套接字的好用处。毕竟,您使用的是异步 IO,这完全解决了可伸缩性问题。
  • @usr,实际上,你是对的,我不需要 Blocking = false 调用,因为我在 Socket 上使用 ReceiveAsync,它忽略了 Blocking 属性。这是我在尝试更简洁的异步方法之前最初实现 Socket 操作的方式的遗留物。
  • @DanBryant 啊,我明白了。我什至不熟悉这种行为。它让我想起了 15 年前在 unix 上被认为是现代套接字编程的东西。过时的东西。

标签: c# .net asynchronous tpl-dataflow


【解决方案1】:

是的,所以……你在这里使用了 TPL 的力量。如果您订阅 TPL DataFlow 样式,您仍然在自己的 while 循环中在后台 Task 中手动接收来自 BufferBlock 的项目这一事实并不是您想要的“方式”。

您要做的是将ActionBlock 链接到BufferBlock 并在其中进行命令处理/发送。这也是您设置MaxDegreeOfParallelism 以控制您想要处理多少并发命令的块。所以这个设置可能看起来像这样:

// Initialization logic to build up the TPL flow
_pendingCommands = new BufferBlock<Command>();
_commandProcessor = new ActionBlock<Command>(this.ProcessCommand);

_pendingCommands.LinkTo(_commandProcessor);

private Task ProcessCommand(Command command)
{
   var response = await _commandProcessor.RunCommand(command.Content);
   this.Send(command.Client, response);
}

然后,在您的关闭代码中,您需要通过在 _pipelineCommands BufferBlock 上调用 Complete 来表明您已完成将项目添加到管道中,然后在 _commandProcessor ActionBlock 上等待完成以确保所有项目都已通过管道。您可以通过获取块的 Completion 属性返回的 Task 并在其上调用 Wait 来做到这一点:

_pendingCommands.Complete();
_commandProcessor.Completion.Wait();

如果您想获得奖励积分,您甚至可以将命令处理与命令发送分开。这将允许您将这些步骤彼此分开配置。例如,您可能需要限制处理命令的线程数,但希望有更多发送响应。您只需在流程中间引入 TransformBlock 即可:

_pendingCommands = new BufferBlock<Command>();
_commandProcessor = new TransformBlock<Command, Tuple<Client, Response>>(this.ProcessCommand);
_commandSender = new ActionBlock<Tuple<Client, Response>(this.SendResponseToClient));

_pendingCommands.LinkTo(_commandProcessor);
_commandProcessor.LinkTo(_commandSender);

private Task ProcessCommand(Command command)
{
   var response = await _commandProcessor.RunCommand(command.Content);

   return Tuple.Create(command, response);
}

private Task SendResponseToClient(Tuple<Client, Response> clientAndResponse)
{
   this.Send(clientAndResponse.Item1, clientAndResponse.Item2);
}

您可能想使用自己的数据结构而不是 Tuple,这只是为了说明目的,但重点是这正是您想要用来分解管道以便您可以控制的那种结构它的各个方面正是您可能需要的。

【讨论】:

  • 这很有帮助;我对 DataFlow 很陌生,但它看起来很有用。令人惊讶的是,即使在 async/await 时代,网络 I/O 仍然如此尴尬。我继续将其简化为单个 ActionBlock(因为它已经具有缓冲功能),初步的单元测试看起来很有希望。
【解决方案2】:

任务默认为后台,这意味着当应用程序终止时,它们也会立即终止。您应该使用线程而不是任务。然后你可以设置:

Thread.IsBackground = false;

这将防止您的应用程序在工作线程运行时终止。 当然,这需要对上述代码进行一些更改。

更何况你在执行shutdown方法的时候,也可以只等待主线程的任何未完成的任务。

我没有看到更好的解决方案。

【讨论】:

  • 不幸的是,这违背了我想要完成的目的,即允许将任务有效分区到现有线程池线程上,而不是在命令到达时创建新线程。
  • 另外,这是一个流程全局解决方案。如果他有多个独立的 NetworkService 在运行怎么办?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-12-14
  • 2021-10-21
  • 2015-12-22
  • 1970-01-01
相关资源
最近更新 更多