DataflowBlock ITargetSource.AsObservable () не запускает OnNext ()

Aug 28 2020

Я пытаюсь использовать блок потока данных, и мне нужно следить за проходящими элементами для модульного тестирования.

Для того , чтобы сделать это, я использую AsObservable()метод на ISourceBlock<T>из моих TransformBlock<Tinput, T>, так что я могу проверить после выполнения , что каждый блок моего трубопровода породившего ожидаемых значений.

Трубопровод

{
   ...
   var observer = new MyObserver<string>();
   _block  = new TransformManyBlock<string, string>(MyHandler, options);
   _block.LinkTo(_nextBlock);
   _block.AsObservable().Subscribe(observer);
   _block.Post("Test");
   ...
}

MyObserver

public class MyObserver<T> : IObserver<T>
{
    public List<Exception> Errors = new List<Exception>();
    public bool IsComplete = false;
    public List<T> Values = new List<T>();

    public void OnCompleted()
    {
        IsComplete = true;
    }

    public void OnNext(T value)
    {
        Values.Add(value);
    }

    public void OnError(Exception e)
    {
        Errors.Add(e);
    }
}

Так что в основном я подписываю своего наблюдателя на блок преобразования и ожидаю, что каждое проходящее значение будет зарегистрировано в моем списке «значений» наблюдателя.

Но, хотя для параметра IsCompleteустановлено значение true и OnError()исключение успешно зарегистрировано, OnNext()метод никогда не вызывается, если он не является последним блоком конвейера ... Я не могу понять, почему, потому что «следующий блок» успешно связан с этим исходным блоком получить данные, доказывающие, что некоторые данные выходят из блока.

Насколько я понимаю, AsObservableпредполагается , что он должен сообщать все значения, выходящие из блока, а не только значения, которые не были использованы другими связанными блоками ...

Что я делаю неправильно ?

Ответы

2 00110001 Aug 28 2020 at 17:08

Ваши сообщения потребляются , _nextBlockпрежде чем вы получите возможность их прочитать.

Если вы закомментируете эту строку, _block.LinkTo(_nextBlock);она, скорее всего, сработает.

AsObservableединственной целью является просто позволить блок будет потребляться от RX . Это не меняет внутреннюю работу блока для передачи сообщений нескольким целям . Для этого нужен специальный блокBroadcastBlock

Я бы предложил транслировать в другой блок и использовать это дляSubscribe

Миссия BroadcastBlock в жизни - дать возможность всем целям, связанным с блоком, получить копию каждого опубликованного элемента.

var options = new DataflowLinkOptions {PropagateCompletion = true};


var broadcastBlock = new BroadcastBlock<string>(x => x);
var bufferBlock = new BufferBlock<string>();
var actionBlock = new ActionBlock<string>(s => Console.WriteLine("Action " + s));

broadcastBlock.LinkTo(bufferBlock, options);
broadcastBlock.LinkTo(actionBlock, options);

bufferBlock.AsObservable().Subscribe(s => Console.WriteLine("peek " + s));

for (var i = 0; i < 5; i++)
   await broadcastBlock.SendAsync(i.ToString());

broadcastBlock.Complete();
await actionBlock.Completion;

Вывод

peek 0
Action 0
Action 1
Action 2
Action 3
Action 4
peek 1
peek 2
peek 3
peek 4