BatchBlock produces batch with elements sent after TriggerBatch()

Viewed 1656

I have a Dataflow pipeline consisting of several blocks. When elements are flowing through my processing pipeline, I want to group them by field A. To do this I have a BatchBlock with high BoundedCapacity. In it I store my elements until I decide that they should be released. So I invoke TriggerBatch() method.

private void Forward(TStronglyTyped data)
{
    if (ShouldCreateNewGroup(data))
    {
        GroupingBlock.TriggerBatch();
    }

 GroupingBlock.SendAsync(data).Wait(SendTimeout);
}

This is how it looks. The problem is, that the batch produced, sometimes contains the next posted element, which shouldn't be there.

To illustrate:

BatchBlock.InputQueue = {A,A,A}
NextElement = B //we should trigger a Batch!
BatchBlock.TriggerBatch()
BatchBlock.SendAsync(B);

In this point I expect my batch to be {A,A,A}, but it is {A,A,A,B}

Like TriggerBatch() was asynchronous, and SendAsync was in fact executed before the batch was actually made.

How can I solve this? I obviously don't want to put Task.Wait(x) in there (I tried, and it works, but then performance is poor, of course).

2 Answers

Here is a specialized version of Loren Paulsen's CreateConditionalBatchBlock method. This one accepts a Func<TItem, TKey> keySelector argument, and emits a new batch every time an item with different key is received.

public static IPropagatorBlock<TItem, TItem[]> CreateConditionalBatchBlock<TItem, TKey>(
    Func<TItem, TKey> keySelector,
    DataflowBlockOptions dataflowBlockOptions = null,
    int maxBatchSize = DataflowBlockOptions.Unbounded,
    IEqualityComparer<TKey> keyComparer = null)
{
    if (keySelector == null) throw new ArgumentNullException(nameof(keySelector));
    if (maxBatchSize < 1 && maxBatchSize != DataflowBlockOptions.Unbounded)
        throw new ArgumentOutOfRangeException(nameof(maxBatchSize));

    keyComparer = keyComparer ?? EqualityComparer<TKey>.Default;
    var options = new ExecutionDataflowBlockOptions();
    if (dataflowBlockOptions != null)
    {
        options.BoundedCapacity = dataflowBlockOptions.BoundedCapacity;
        options.CancellationToken = dataflowBlockOptions.CancellationToken;
        options.MaxMessagesPerTask = dataflowBlockOptions.MaxMessagesPerTask;
        options.TaskScheduler = dataflowBlockOptions.TaskScheduler;
    }

    var output = new BufferBlock<TItem[]>(options);

    var queue = new Queue<TItem>(); // Synchronization is not needed
    TKey previousKey = default;

    var input = new ActionBlock<TItem>(async item =>
    {
        var key = keySelector(item);
        if (queue.Count > 0 && !keyComparer.Equals(key, previousKey))
        {
            await output.SendAsync(queue.ToArray()).ConfigureAwait(false);
            queue.Clear();
        }
        queue.Enqueue(item);
        previousKey = key;

        if (queue.Count == maxBatchSize)
        {
            await output.SendAsync(queue.ToArray()).ConfigureAwait(false);
            queue.Clear();
        }
    }, options);

    _ = input.Completion.ContinueWith(async t =>
    {
        if (queue.Count > 0)
        {
            await output.SendAsync(queue.ToArray()).ConfigureAwait(false);
            queue.Clear();
        }
        if (t.IsFaulted)
        {
            ((IDataflowBlock)output).Fault(t.Exception.InnerException);
        }
        else
        {
            output.Complete();
        }
    }, TaskScheduler.Default);

    return DataflowBlock.Encapsulate(input, output);
}
Related