【问题标题】:How to restore the order of a shuffled Dataflow pipeline?如何恢复打乱的 Dataflow 管道的顺序?
【发布时间】:2020-12-18 14:08:53
【问题描述】:

我有一个数据流管道,它由多个处理异构文档(XLS、PDF 等)的块组成。每种类型的文档都由专用的TransformBlock 处理。在管道的末端,我有一个ActionBlock,它接收所有已处理的文档,并将它们一一上传到 Web 服务器。我的问题是我找不到一种方法来满足按照最初输入管道的相同顺序上传文档的要求。例如,我不能使用 EnsureOrdered 选项来发挥我的优势,因为此选项配置单个块的行为,而不是并行工作的多个块的行为。我的要求是:

  1. 按特定顺序将文档插入管道中。
  2. 根据文档类型以不同方式处理每个文档。
  3. 应按顺序处理特定类型的文档。
  4. 可以(并且应该)并行处理不同类型的文档。
  5. 所有文件都应在处理完毕后尽快上传。
  6. 文档必须按顺序上传,并按照它们在管道中输入的顺序。

例如要求文档#8必须在文档#7之后上传,即使它是在文档#7之前处理的。

第五个要求的意思是我等不及所有文档都处理完,然后按索引排序,最后上传。上传必须与处理同时进行。

这是我正在尝试做的一个最小示例。为简单起见,我没有使用IDocument 接口的实例来提供块,而是使用简单的整数。每个整数的值代表它进入管道的顺序,以及必须上传的顺序:

var xlsBlock = new TransformBlock<int, int>(document =>
{
    int duration = 300 + document % 3 * 300;
    Thread.Sleep(duration); // Simulate CPU-bound work
    return document;
});
var pdfBlock = new TransformBlock<int, int>(document =>
{
    int duration = 100 + document % 5 * 200;
    Thread.Sleep(duration); // Simulate CPU-bound work
    return document;
});

var uploader = new ActionBlock<int>(async document =>
{
    Console.WriteLine($"Uploading document #{document}");
    await Task.Delay(500); // Simulate I/O-bound work
});

xlsBlock.LinkTo(uploader);
pdfBlock.LinkTo(uploader);

foreach (var document in Enumerable.Range(1, 10))
{
    if (document % 2 == 0)
        xlsBlock.Post(document);
    else
        pdfBlock.Post(document);
}
xlsBlock.Complete();
pdfBlock.Complete();
_ = Task.WhenAll(xlsBlock.Completion, pdfBlock.Completion)
    .ContinueWith(_ => uploader.Complete());

await uploader.Completion;

输出是:

Uploading document #1
Uploading document #2
Uploading document #3
Uploading document #5
Uploading document #4
Uploading document #7
Uploading document #6
Uploading document #9
Uploading document #8
Uploading document #10

(Try it on Fiddle)

理想的顺序是#1、#2、#3、#4、#5、#6、#7、#8、#9、#10。

在将已处理文档发送到uploader 块之前,如何恢复已处理文档的顺序?

澄清:通过将多个特定的TransformBlocks 替换为单个通用的TransformBlock 来彻底改变管道的架构不是一种选择。理想的情况是拦截处理器和上传者之间的单个块,这将恢复文档的顺序。

【问题讨论】:

  • 正常的方法是在传输的开头添加一个序列号,以便块可以按正确的顺序重新组装。
  • @jdweng 您可以假设序列号已经是document 对象的属性。我可以向它们添加public long SequenceNumber 属性,并正确初始化它。问题是如何在它们被所有这些不同的块处理后重新组装它们。
  • 上传者需要在每个块前添加一个序列号,以便服务器收到块时,服务器可以按正确的顺序组合。您不能使用在处理过程中被删除的号码。
  • @jdweng 我不能将订单恢复委托给网络服务器(如果你是这个意思)。我必须在我自己的程序中这样做。
  • 那你必须在上传前做,不能并行上传。

标签: c# task-parallel-library tpl-dataflow


【解决方案1】:

uploader 应该将文档添加到已完成文档的排序列表中,并检查添加的文档是否是下一个应该上传的文档。如果它应该从排序列表中删除并上传所有文档,直到缺少一个。

还有一个同步问题。对该排序列表的访问必须跨线程同步。但是您希望所有线程都在做某事,而不是等待其他线程完成它们的工作。所以,uploader 应该与这样的列表一起使用:

  • 在同步锁内将新文档添加到列表中,并释放锁
  • 在循环中
    • 再次输入相同的同步锁,
    • 如果设置了upload_in_progress 标志,则什么也不做并返回。
    • 检查是否应上传列表顶部的文档,
      • 如果不是,则重置upload_in_progress 标志,然后返回。
      • 否则从列表中删除文档,
      • 设置upload_in_progress标志,
      • 释放锁,
      • 上传文档。

我希望我想象的没错。正如你所看到的,让它既安全又高效是很棘手的。在大多数情况下,肯定有一种方法可以只使用一个锁,但它不会增加太多效率。 upload_in_progress 标志在任务之间共享,就像列表本身一样。

【讨论】:

  • 感谢 Dialecticus 的回答!我认为你的想法倾向于一个可能的解决方案。你认为我可以使用SortedList&lt;TKey, TValue&gt; 类来解决这个问题吗?
  • 确实如此。 TKey 是文档序号,TValue 是文档本身。顺便说一句,对列表的访问(以及上传本身)必须受到某种同步机制的保护,例如lock。现在实现起来很棘手,但对于 Stack Overflow 来说,这也是一个很好的新问题,一旦你可以看到并正式化这个问题。
  • 好的,我会尝试使用SortedList 来实现你的想法,然后看看它的发展方向。
  • 我为同步问题添加了一个可能的解决方案。祝你好运。
  • 谢谢!顺便说一句,uploader 按顺序(一次一个)上传文件,所以我认为我可以摆脱同步。但我会记住这一点,以防将来需求发生变化。
【解决方案2】:

我设法实现了一个数据流块,它可以根据包含已处理文档的排序列表的 Dialecticus idea 恢复我的混洗管道的顺序。我最终使用了一个简单的Dictionary,而不是SortedList,这似乎也同样有效。

/// <summary>Creates a dataflow block that restores the order of
/// a shuffled pipeline.</summary>
public static IPropagatorBlock<T, T> CreateRestoreOrderBlock<T>(
    Func<T, long> indexSelector,
    long startingIndex = 0L,
    DataflowBlockOptions options = null)
{
    if (indexSelector == null) throw new ArgumentNullException(nameof(indexSelector));
    var executionOptions = new ExecutionDataflowBlockOptions();
    if (options != null)
    {
        executionOptions.CancellationToken = options.CancellationToken;
        executionOptions.BoundedCapacity = options.BoundedCapacity;
        executionOptions.EnsureOrdered = options.EnsureOrdered;
        executionOptions.TaskScheduler = options.TaskScheduler;
        executionOptions.MaxMessagesPerTask = options.MaxMessagesPerTask;
        executionOptions.NameFormat = options.NameFormat;
    }

    var buffer = new Dictionary<long, T>();
    long minIndex = startingIndex;

    IEnumerable<T> Transform(T item)
    {
        // No synchronization needed because MaxDegreeOfParallelism = 1
        long index = indexSelector(item);
        if (index < startingIndex)
            throw new InvalidOperationException($"Index {index} is out of range.");
        if (index < minIndex)
            throw new InvalidOperationException($"Index {index} has been consumed.");
        if (!buffer.TryAdd(index, item)) // .NET Core only API
            throw new InvalidOperationException($"Index {index} is not unique.");
        while (buffer.Remove(minIndex, out var minItem)) // .NET Core only API
        {
            minIndex++;
            yield return minItem;
        }
    }

    // Ideally the assertion buffer.Count == 0 should be checked on the completion
    // of the block.
    return new TransformManyBlock<T, T>(Transform, executionOptions);
}

使用示例:

var xlsBlock = new TransformBlock<int, int>(document =>
{
    int duration = 300 + document % 3 * 300;
    Thread.Sleep(duration); // Simulate CPU-bound work
    return document;
});
var pdfBlock = new TransformBlock<int, int>(document =>
{
    int duration = 100 + document % 5 * 200;
    Thread.Sleep(duration); // Simulate CPU-bound work
    return document;
});

var orderRestorer = CreateRestoreOrderBlock<int>(
    indexSelector: document => document, startingIndex: 1L);

var uploader = new ActionBlock<int>(async document =>
{
    Console.WriteLine($"{DateTime.Now:HH:mm:ss.fff} Uploading document #{document}");
    await Task.Delay(500); // Simulate I/O-bound work
});

xlsBlock.LinkTo(orderRestorer);
pdfBlock.LinkTo(orderRestorer);
orderRestorer.LinkTo(uploader, new DataflowLinkOptions { PropagateCompletion = true });

foreach (var document in Enumerable.Range(1, 10))
{
    if (document % 2 == 0)
        xlsBlock.Post(document);
    else
        pdfBlock.Post(document);
}
xlsBlock.Complete();
pdfBlock.Complete();
_ = Task.WhenAll(xlsBlock.Completion, pdfBlock.Completion)
    .ContinueWith(_ => orderRestorer.Complete());

await uploader.Completion;

输出:

09:24:18.846 Uploading document #1
09:24:19.436 Uploading document #2
09:24:19.936 Uploading document #3
09:24:20.441 Uploading document #4
09:24:20.942 Uploading document #5
09:24:21.442 Uploading document #6
09:24:21.941 Uploading document #7
09:24:22.441 Uploading document #8
09:24:22.942 Uploading document #9
09:24:23.442 Uploading document #10

Try it on Fiddle,具有 .NET Framework 兼容版本)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-17
    • 2016-06-20
    • 2021-06-09
    相关资源
    最近更新 更多