DataflowBlock ITargetSource.AsObservable () löst OnNext () nicht aus

Aug 28 2020

Ich versuche, einen Datenflussblock zu verwenden, und ich muss die Elemente ausspionieren, die für Unit-Tests durchlaufen werden.

Zu diesem Zweck verwende ich die AsObservable()Methode on ISourceBlock<T>von my TransformBlock<Tinput, T>, damit ich nach der Ausführung überprüfen kann, ob jeder Block meiner Pipeline die erwarteten Werte generiert hat.

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

Im Grunde genommen abonniere ich meinen Beobachter für den Transformationsblock und erwarte, dass jeder durchlaufende Wert in meiner Beobachterliste "Werte" registriert wird.

Aber während das IsCompleteauf true gesetzt ist und die OnError()Ausnahme erfolgreich registriert wird, wird die OnNext()Methode nie aufgerufen, es sei denn, es ist der letzte Block der Pipeline ... Ich kann nicht herausfinden, warum, weil der "nächste Block" erfolgreich mit diesem sourceBlock verknüpft ist Empfangen Sie die Daten, um zu beweisen, dass einige Daten den Block verlassen.

Soweit ich weiß, AsObservablesoll das alle Werte melden, die den Block verlassen, und nicht nur die Werte, die nicht von anderen verknüpften Blöcken verbraucht wurden ...

Was mache ich falsch ?

Antworten

2 00110001 Aug 28 2020 at 17:08

Ihre Nachrichten werden verbraucht durch , _nextBlockbevor Sie eine Chance bekommen , sie zu lesen.

Wenn Sie diese Zeile _block.LinkTo(_nextBlock);auskommentieren, würde es wahrscheinlich funktionieren.

AsObservablealleiniger Zweck es ist nur zu ermöglichen , Block wird verbraucht von RX . Die interne Arbeitsweise des Blocks wird nicht geändert, um Nachrichten an mehrere Ziele zu senden . Dafür benötigen Sie einen speziellen BlockBroadcastBlock

Ich würde vorschlagen , in einen anderen Block zu senden und diesen zu verwendenSubscribe

BroadcastBlocks Lebensaufgabe ist es, allen vom Block verknüpften Zielen zu ermöglichen, eine Kopie jedes veröffentlichten Elements zu erhalten

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;

Ausgabe

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