DataflowBlock ITargetSource.AsObservable () não dispara OnNext ()

Aug 28 2020

Estou tentando usar um bloco de fluxo de dados e preciso espionar os itens que estão passando para o teste de unidade.

Para fazer isso, estou usando o AsObservable()método on ISourceBlock<T>do meu TransformBlock<Tinput, T>, para que possa verificar após a execução se cada bloco do meu pipeline gerou os valores esperados.

Pipeline

{
   ...
   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);
    }
}

Então, basicamente, eu inscrevo meu observador no transformblock e espero que cada valor passando seja registrado na minha lista de "valores" do observador.

Mas, enquanto o IsCompleteestá definido como verdadeiro e OnError()registra a exceção com sucesso, o OnNext()método nunca é chamado a menos que seja o último bloco do pipeline ... Não consigo descobrir o porquê, porque o "nextblock" vinculou a este sourceBlock com sucesso receber os dados, provando que alguns dados estão saindo do bloco.

Pelo que entendi, o AsObservabledeve reportar todos os valores que saem do bloco e não apenas os valores que não foram consumidos por outros blocos vinculados ...

O que estou fazendo errado ?

Respostas

2 00110001 Aug 28 2020 at 17:08

Suas mensagens estão sendo consumidos por _nextBlockantes de ter a chance de lê-los.

Se você comentar esta linha _block.LinkTo(_nextBlock);, provavelmente funcionará.

AsObservableO único propósito é permitir que um bloco seja consumido do RX . Isso não altera o funcionamento interno do bloco para transmitir mensagens para vários destinos . Você precisa de um bloco especial para issoBroadcastBlock

Eu sugeriria transmitir para outro bloco e usá-lo paraSubscribe

A missão do BroadcastBlock na vida é permitir que todos os alvos vinculados ao bloco tenham uma cópia de cada elemento publicado

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;

Resultado

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