DataflowBlock ITargetSource.AsObservable () ne déclenchant pas OnNext ()

Aug 28 2020

J'essaie d'utiliser un dataflowblock et j'ai besoin d'espionner les éléments qui passent pour les tests unitaires.

Pour ce faire, j'utilise la AsObservable()méthode on ISourceBlock<T>of my TransformBlock<Tinput, T>, donc je peux vérifier après exécution que chaque bloc de mon pipeline a généré les valeurs attendues.

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

Donc, fondamentalement, j'abonne mon observateur au transformblock, et je m'attends à ce que chaque valeur qui passe soit enregistrée dans ma liste de "valeurs" d'observateur.

Mais, alors que le IsCompleteest défini sur true et que l' OnError()exception d'inscription a réussi, la OnNext()méthode n'est jamais appelée à moins que ce ne soit le dernier bloc du pipeline ... Je ne peux pas comprendre pourquoi, car le "nextblock" lié à ce sourceBlock avec succès recevoir les données, prouvant que certaines données sortent du bloc.

D'après ce que je comprends, le AsObservableest censé signaler toutes les valeurs sortant du bloc et pas seulement les valeurs qui n'ont pas été consommées par d'autres blocs liés ...

Qu'est-ce que je fais mal ?

Réponses

2 00110001 Aug 28 2020 at 17:08

Vos messages sont consommés par _nextBlockavant d'obtenir une chance de les lire.

Si vous commentez cette ligne, _block.LinkTo(_nextBlock);cela fonctionnerait probablement.

AsObservablele seul but est simplement de permettre à un bloc d'être consommé à partir de RX . Cela ne change pas le fonctionnement interne du bloc pour diffuser des messages à plusieurs cibles . Vous avez besoin d'un bloc spécial pour celaBroadcastBlock

Je suggérerais de diffuser vers un autre bloc et de l'utiliser pourSubscribe

La mission de BroadcastBlock dans la vie est de permettre à toutes les cibles liées à partir du bloc d'obtenir une copie de chaque élément publié

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;

Production

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