DataflowBlock ITargetSource.AsObservable () OnNext () 'i tetiklemiyor
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
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