DataflowBlock ITargetSourceAsObservable () ไม่ทริกเกอร์ OnNext ()
ฉันกำลังพยายามใช้ dataflowblock และฉันต้องการสอดแนมรายการที่ผ่านเพื่อทดสอบหน่วย
ในการดำเนินการนี้ฉันใช้AsObservable()วิธีการISourceBlock<T>ของฉันTransformBlock<Tinput, T>ดังนั้นฉันจึงสามารถตรวจสอบหลังการดำเนินการได้ว่าแต่ละบล็อกของไปป์ไลน์ของฉันสร้างค่าที่คาดไว้
ท่อ
{
...
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);
}
}
โดยพื้นฐานแล้วฉันสมัครผู้สังเกตการณ์ของฉันกับ transformblock และฉันคาดหวังว่าแต่ละค่าที่ผ่านไปจะได้รับการลงทะเบียนในรายการ "ค่า" ผู้สังเกตการณ์ของฉัน
แต่ในขณะที่IsCompleteตั้งค่าเป็นจริงและOnError()ข้อยกเว้นในการลงทะเบียนสำเร็จOnNext()เมธอดจะไม่ถูกเรียกเว้นแต่ว่าจะเป็นบล็อกสุดท้ายของไปป์ไลน์ ... ฉันไม่สามารถหาสาเหตุได้เพราะ "nextblock" ที่เชื่อมโยงกับ sourceBlock นี้สำเร็จ รับข้อมูลโดยพิสูจน์ว่าข้อมูลบางส่วนกำลังออกจากบล็อก
จากสิ่งที่ฉันเข้าใจAsObservableควรจะรายงานทุกค่าที่ออกจากบล็อกและไม่เพียง แต่ค่าที่บล็อกอื่นที่เชื่อมโยงไม่ได้ใช้เท่านั้น
ผมทำอะไรผิดหรือเปล่า ?
คำตอบ
ข้อความของคุณจะถูกบริโภคโดย_nextBlockก่อนที่คุณจะได้รับโอกาสที่จะอ่านพวกเขา
หากคุณแสดงความคิดเห็นในบรรทัดนี้_block.LinkTo(_nextBlock);ก็น่าจะใช้ได้
AsObservableวัตถุประสงค์เพียงอย่างเดียวเป็นเพียงเพื่อให้บล็อกที่จะบริโภคจากRX มันไม่ได้เปลี่ยนการทำงานภายในของบล็อกที่จะออกอากาศข้อความไปยังหลายเป้าหมาย คุณต้องมีบล็อกพิเศษสำหรับสิ่งนั้นBroadcastBlock
ฉันขอแนะนำให้ออกอากาศไปยังบล็อกอื่นและใช้เพื่อSubscribe
ภารกิจของ BroadcastBlock ในชีวิตคือการเปิดใช้งานเป้าหมายทั้งหมดที่เชื่อมโยงจากบล็อกเพื่อรับสำเนาของทุกองค์ประกอบที่เผยแพร่
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;
เอาต์พุต
peek 0
Action 0
Action 1
Action 2
Action 3
Action 4
peek 1
peek 2
peek 3
peek 4