DataflowBlock ITargetSource.AsObservable () não dispara OnNext ()
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
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