DataflowBlock ITargetSource.AsObservable () tidak memicu OnNext ()

Aug 28 2020

Saya mencoba menggunakan dataflowblock dan saya perlu memata-matai item yang melewati pengujian unit.

Untuk melakukan ini, saya menggunakan AsObservable()metode on ISourceBlock<T>of my TransformBlock<Tinput, T>, jadi saya dapat memeriksa setelah eksekusi bahwa setiap blok pipeline saya telah menghasilkan nilai yang diharapkan.

Pipa

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

Jadi pada dasarnya saya mendaftarkan pengamat saya ke transformblock, dan saya berharap bahwa setiap nilai yang lewat terdaftar dalam daftar "nilai" pengamat saya.

Tapi, sementara IsCompletedisetel ke true, dan OnError()pengecualian register berhasil, OnNext()metode tidak pernah dipanggil kecuali itu adalah blok terakhir dari pipeline ... Saya tidak tahu mengapa, karena "blok berikutnya" terhubung ke sourceBlock ini dengan sukses menerima data, membuktikan bahwa beberapa data keluar dari blok.

Dari apa yang saya pahami, AsObservableseharusnya melaporkan setiap nilai yang keluar dari blok dan tidak hanya nilai yang belum dikonsumsi oleh blok terkait lainnya ...

Apa yang saya lakukan salah?

Jawaban

2 00110001 Aug 28 2020 at 17:08

Pesan Anda sedang dikonsumsi oleh _nextBlocksebelum Anda mendapatkan kesempatan untuk membacanya.

Jika Anda mengomentari baris _block.LinkTo(_nextBlock);ini, kemungkinan besar akan berhasil.

AsObservabletujuan utamanya adalah hanya untuk memungkinkan blok untuk dikonsumsi dari RX . Itu tidak mengubah kerja internal blok untuk menyiarkan pesan ke beberapa target . Anda membutuhkan blok khusus untuk ituBroadcastBlock

Saya akan menyarankan penyiaran ke blok lain dan menggunakannya untukSubscribe

Misi BroadcastBlock dalam hidup adalah mengaktifkan semua target yang ditautkan dari blok untuk mendapatkan salinan dari setiap elemen yang diterbitkan

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;

Keluaran

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