您真的需要 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);
}
}