DataflowBlock ITargetSource.AsObservable () OnNext () 'i tetiklemiyor

Aug 28 2020

Bir veri akışı bloğu kullanmaya çalışıyorum ve birim testi için geçen öğeleri gözetlemem gerekiyor.

Bunu yapmak için, ben kullanıyorum AsObservable()üzerinde yöntemini ISourceBlock<T>Sesimin TransformBlock<Tinput, T>benim boru hattının her blok beklenen değerleri yarattı yürütme işleminden sonra kontrol edebilirsiniz.

Boru hattı

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

Yani temelde gözlemcimi transformblock'a abone oluyorum ve geçen her değerin gözlemci "değerler" listeme kaydedilmesini bekliyorum.

Ancak, IsCompletetrue olarak ayarlıyken ve OnError()başarılı kayıt istisnası olsa da, OnNext()yöntem, ardışık düzenin son bloğu olmadığı sürece asla çağrılmaz ... Nedenini anlayamıyorum, çünkü bu sourceBlock'a başarıyla bağlanan "sonraki blok" bazı verilerin bloktan çıktığını kanıtlayan verileri alır.

Anladığım kadarıyla, AsObservablesadece diğer bağlantılı bloklar tarafından tüketilmeyen değerleri değil, bloktan çıkan her değeri rapor etmesi gerekiyor ...

Neyi yanlış yapıyorum ?

Yanıtlar

2 00110001 Aug 28 2020 at 17:08

Mesajlarınız ediliyor tüketilen tarafından _nextBlockbunları okumak için bir şans elde önce.

Bu satırı yorumlarsanız _block.LinkTo(_nextBlock);muhtemelen işe yarayacaktır.

AsObservableTek amacı sadece izin vermektir blok olmak tüketilen gelen RX . Mesajları birden fazla hedefe yayınlamak için bloğun dahili çalışmasını değiştirmez . Bunun için özel bir bloğa ihtiyacın varBroadcastBlock

Öneririm yayın diğerine bloğu ve buna kullanılarakSubscribe

BroadcastBlock'un hayattaki misyonu, bloktan bağlantılı tüm hedeflerin yayınlanan her öğenin bir kopyasını almasını sağlamaktır.

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;

Çıktı

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