【问题标题】:Is there a way to only process one task at a time per linked block in TPL Dataflow?有没有办法在 TPL 数据流中的每个链接块一次只处理一个任务?
【发布时间】:2018-05-14 18:16:35
【问题描述】:

我有两个类,serviceclient 和一个服务。 serviceclient 将生成将由服务处理的消息。 ServiceClient 消息应该按照严格的 FIFO 顺序处理,下一条消息仅在前一条消息完成处理后才可用。

为了解决这个问题,我在每个服务客户端中放置了一个操作块,它直接调用服务来处理客户端消息,我认为这可以正常工作,但需要额外的依赖注入。我想知道是否有一种方法可以设置它,以便服务可以直接链接到 serviceclient 消息块,以便它可以同时处理来自多个 serviceclient 消息块的消息,但一次只能处理来自任何特定 serviceclient 的一条消息?

一些具有所需功能的代码::

    static void Main(string[] args)
    {
        var service = new Service();

        //I would like messages from these clients to be processed concurrently by service, but only one at a time per client. 
        //So if Client A has two messages in queue(a1, a2) and B has 3(b1,b2,b3), it will immediately take a1&b1. If a1 finishes, it will then take a2. if b1 finishes it will take b2, same with b2 and b3. It would never process a1 concurrently with a2, or b1 concurrently with b2 or b3. 
        service.AddClient(new ServiceClient());
        service.AddClient(new ServiceClient());
    }

    interface IServiceMessage
    {
        string Message { get; }
    }

    class ServiceClient
    {
        public BufferBlock<IServiceMessage> clientServiceMsgs = new BufferBlock<IServiceMessage>();

        public ServiceClient() {
            //run task to populate bufferblock 
        }
    }

    class Service
    {
        ActionBlock<IServiceMessage> processServiceMsgsBlock;

        public Service() {
            processServiceMsgsBlock = new ActionBlock<IServiceMessage>(ProcessServiceMessage);
        }  

        public async Task ProcessServiceMessage(IServiceMessage msg) {
            //process stuff
            return;
        }

        public void AddClient(ServiceClient client)
        {
            client.clientServiceMsgs.LinkTo(processServiceMsgsBlock);
        }
    }

【问题讨论】:

  • 一些示例代码说明了您尝试做的事情,这将大大有助于我们为您提供帮助。
  • 让我知道这是否澄清了事情

标签: c# tpl-dataflow


【解决方案1】:

您真的需要 TPL DataFlow 吗?我会给你两个不带的选项,最后一个是你问的 TPL 数据流。

每个选项共享以下内容。

public class ServiceClient
{
    public async Task<IServiceMessage> GetAsync(string filter)
    {
        await Task.Yield();
        return new ServiceMessage() { Message = Guid.NewGuid().ToString() };
    }
}

public class Service
{
    public async Task ProcessAsync(IServiceMessage message)
    {
        // do something with it
        await Task.Delay(10);
    }
}

public interface IServiceMessage
{
    string Message { get; }
}

public class ServiceMessage : IServiceMessage
{
    public string Message { get; set; }
}

选项 1 - 服务启动读取和写入数据的任务。只需将 Service 和 ServiceClient 耦合到第三个服务中即可。

class Program
{
    static void Main(string[] args)
    {
        var cts = new CancellationTokenSource();
        var service = new Service();
        var serviceClient = new ServiceClient();
        var processor = new ProducerConsumerService(serviceClient, service);
        processor.Process("A", cts.Token);
        processor.Process("B", cts.Token);
        processor.Process("C", cts.Token);
        processor.Process("D", cts.Token);

        Console.WriteLine("Press any key to shutdown");
        Console.Read();

        cts.Cancel();
        processor.WaitForCompletion();
    }
}

public class ProducerConsumerService
{
    private List<Task> _processTasks;
    private ServiceClient _serviceClient;
    private Service _service;

    public ProducerConsumerService(ServiceClient serviceClient, Service service)
    {
        _serviceClient = serviceClient;
        _service = service;

        _processTasks = new List<Task>();
    }

    public void Process(string filter, CancellationToken token)
    {
        _processTasks.Add(Task.Run(() =>
        {
            while (!token.IsCancellationRequested)
            {
                var message = _serviceClient.Get(filter);
                _service.Process(message);
            }
        }));
    }

    public void WaitForCompletion()
    {
        Task.WaitAll(_processTasks.ToArray(), TimeSpan.FromSeconds(10));
    }
}

选项 2 - 与选项 1 相同,但有两个任务和一个 BlockingCollection,它在生产者 (ServiceClient) 和消费者 (Service) 之间提供有界缓冲区。

public class ProducerBufferConsumerService
{
    private List<Task> _producerTasks;
    private List<Task> _consumerTasks;
    private ServiceClient _serviceClient;
    private Service _service;

    public ProducerBufferConsumerService(ServiceClient serviceClient, Service service)
    {
        _serviceClient = serviceClient;
        _service = service;

        _producerTasks = new List<Task>();
        _consumerTasks = new List<Task>();
    }

    public void Process(CancellationToken token)
    {
        var buffer = new BlockingCollection<IServiceMessage>(1000);
        _producerTasks.Add(Task.Run(async () =>
        {
            while (!token.IsCancellationRequested)
            {
                var message = await _serviceClient.GetAsync();
                buffer.Add(message, token);
            }

            buffer.CompleteAdding();
        }));

        _consumerTasks.Add(Task.Run(async () =>
        {
            while (!token.IsCancellationRequested && !buffer.IsAddingCompleted)
            {
                var message = buffer.Take(token);
                await _service.ProcessAsync(message);
            }
        }));
    }

    public void WaitForCompletion()
    {
        Task.WaitAll(_producerTasks.ToArray(), 10000);
        Task.WaitAll(_consumerTasks.ToArray(), 10000);
    }
}

选项 3 - 并行处理,同时保持给定 ServiceClient 的顺序

不完全符合您的要求,但很接近。您可以并行化传入的 ServiceClient 数据流,同时保持正确的顺序。在这个例子中,所有的 ServiceClients 发布到一个单独的块,然后馈送到 9 个分区的 actionblocks..

为给定服务客户端的每条消息提供一个 Guid Id。然后使用 GetHashCode() 将其转换为数字。将该数字减少到 1-9 的范围。创建 9 个 ActionBlock 并使用 LinkTo 方法中的 lambda 将链接限制为该范围内的单个数字。这样,我们创建了 9 个分区,每个分区处理数据的一个子范围。您可以通过独立操作安全地以并行方式进行处理,同时保持给定 id 的 FIFO。

class Program
{
    static void Main(string[] args)
    {
        var cts = new CancellationTokenSource();
        var filters = new List<string>() { "A", "B", "C", "D" };
        var service = new Service();
        var serviceClient = new ServiceClient();
        var partitioningService = new PartitioningService(serviceClient, service);
        var processingTask = Task.Run(() => partitioningService.Process(filters, cts.Token));


        Console.WriteLine("Press any key to shutdown");
        Console.ReadKey();

        cts.Cancel();
        processingTask.Wait(10000);
    }
}

public interface IServiceMessage
{
    string Message { get; }
    Guid Id { get; set; }
}

public class ServiceMessage : IServiceMessage
{
    public string Message { get; set; }
    public Guid Id { get; set;  }
}

public class RoutedMessage
{
    public IServiceMessage Message { get; set; }
    public int PartitionId { get; set; }
}

public class PartitioningService
{
    private ServiceClient _serviceClient;
    private Service _service;

    public PartitioningService(ServiceClient serviceClient, Service service)
    {
        _serviceClient = serviceClient;
        _service = service;
    }

    public void Process(List<string> filters, CancellationToken token)
    {
        var linkOptions = new DataflowLinkOptions { PropagateCompletion = true };
        Func<IServiceMessage, RoutedMessage> partitioner = x => new RoutedMessage
            { 
                Message = x,
                PartitionId = x.Id.GetHashCode() / 1000000000
        };

        var partitionerBlock = new TransformBlock<IServiceMessage, RoutedMessage>(partitioner);
        var actionBlock1 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock2 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock3 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock4 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock5 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock6 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock7 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock8 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));
        var actionBlock9 = new ActionBlock<RoutedMessage>(async (RoutedMessage msg) => await _service.ProcessAsync(msg));

        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == -4);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == -3);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == -2);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == -1);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == 0);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == 1);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == 2);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == 3);
        partitionerBlock.LinkTo(actionBlock1, linkOptions, msg => msg.PartitionId == 4);

        var tasks = new List<Task>();
        foreach (var filter in filters)
        {
            tasks.Add(Task.Run(async () =>
                {
                    Guid filterId = Guid.NewGuid();

                    while (!token.IsCancellationRequested)
                    {
                        var message = await _serviceClient.GetAsync(filter);
                        message.Id = filterId;
                        await partitionerBlock.SendAsync(message);
                    }
                }));
        }

        while (!token.IsCancellationRequested)
            Thread.Sleep(100);

        partitionerBlock.Complete();
        actionBlock1.Completion.Wait();
        actionBlock2.Completion.Wait();
        actionBlock3.Completion.Wait();
        actionBlock4.Completion.Wait();
        actionBlock5.Completion.Wait();
        actionBlock6.Completion.Wait();
        actionBlock7.Completion.Wait();
        actionBlock8.Completion.Wait();
        actionBlock9.Completion.Wait();

        Task.WaitAll(tasks.ToArray(), 10000);
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-04-09
    • 1970-01-01
    • 1970-01-01
    • 2022-01-05
    • 1970-01-01
    • 2020-05-17
    • 1970-01-01
    相关资源
    最近更新 更多