DataflowBlock ITargetSourceAsObservable () ไม่ทริกเกอร์ OnNext ()

Aug 28 2020

ฉันกำลังพยายามใช้ 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ควรจะรายงานทุกค่าที่ออกจากบล็อกและไม่เพียง แต่ค่าที่บล็อกอื่นที่เชื่อมโยงไม่ได้ใช้เท่านั้น

ผมทำอะไรผิดหรือเปล่า ?

คำตอบ

2 00110001 Aug 28 2020 at 17:08

ข้อความของคุณจะถูกบริโภคโดย_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