DataflowBlock ITargetSource.AsObservable () no activa OnNext ()

Aug 28 2020

Estoy tratando de usar un bloque de flujo de datos y necesito espiar los elementos que pasan para realizar pruebas unitarias.

Para hacer esto, estoy usando el AsObservable()método ISourceBlock<T>de mi TransformBlock<Tinput, T>, por lo que puedo verificar después de la ejecución que cada bloque de mi canalización haya generado los valores esperados.

Tubería

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

Así que básicamente suscribo a mi observador al bloque de transformación, y espero que cada valor que pasa se registre en mi lista de "valores" de observador.

Pero, mientras que IsCompletese establece en verdadero, y la OnError()excepción de registro con éxito, el OnNext()método nunca se llama a menos que sea el último bloque de la canalización ... No puedo entender por qué, porque el "nextblock" vinculado a este sourceBlock correctamente recibir los datos, lo que demuestra que algunos datos están saliendo del bloque.

Por lo que tengo entendido, AsObservablese supone que debe informar todos los valores que salen del bloque y no solo los valores que no han sido consumidos por otros bloques vinculados ...

Qué estoy haciendo mal ?

Respuestas

2 00110001 Aug 28 2020 at 17:08

Sus mensajes están siendo consumidos por _nextBlockantes de recibir una oportunidad de leerlos.

Si comenta esta línea _block.LinkTo(_nextBlock);, probablemente funcionará.

AsObservableEl único propósito es permitir que se consuma un bloque de RX . No cambia el funcionamiento interno del bloque para transmitir mensajes a múltiples objetivos . Necesitas un bloque especial para esoBroadcastBlock

Sugeriría transmitir a otro bloque y usarlo paraSubscribe

La misión de BroadcastBlock en la vida es permitir que todos los objetivos vinculados desde el bloque obtengan una copia 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;

Salida

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