DataflowBlock ITargetSource.AsObservable () không kích hoạt OnNext ()

Aug 28 2020

Tôi đang cố gắng sử dụng dataflowblock và tôi cần theo dõi các mục đi qua để kiểm tra đơn vị.

Để thực hiện việc này, tôi đang sử dụng AsObservable()phương thức trên ISourceBlock<T>của tôi TransformBlock<Tinput, T>, vì vậy tôi có thể kiểm tra sau khi thực hiện rằng mỗi khối trong đường ống của tôi đã tạo ra các giá trị mong đợi.

Đường ống

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

Vì vậy, về cơ bản tôi đăng ký người quan sát của mình vào khối chuyển đổi và tôi hy vọng rằng mỗi giá trị đi qua sẽ được đăng ký trong danh sách "giá trị" người quan sát của tôi.

Tuy nhiên, trong khi IsCompleteđược đặt thành true và OnError()đăng ký thành công ngoại lệ, OnNext()phương thức không bao giờ được gọi trừ khi nó là khối cuối cùng của đường ống ... Tôi không thể tìm ra lý do tại sao, bởi vì "nextblock" được liên kết với sourceBlock này thành công nhận dữ liệu, chứng minh rằng một số dữ liệu đang thoát ra khỏi khối.

Theo những gì tôi hiểu, AsObservablenghĩa là phải báo cáo mọi giá trị thoát ra khỏi khối và không chỉ các giá trị chưa được sử dụng bởi các khối được liên kết khác ...

Tôi đang làm gì sai?

Trả lời

2 00110001 Aug 28 2020 at 17:08

Tin nhắn của bạn đang được tiêu thụ bởi _nextBlocktrước khi bạn có cơ hội để đọc chúng.

Nếu bạn nhận xét ra dòng này, _block.LinkTo(_nextBlock);nó có thể sẽ hoạt động.

AsObservablemục đích duy nhất là chỉ cho phép một khối được tiêu thụ từ RX . Nó không thay đổi hoạt động bên trong của khối để phát thông báo đến nhiều mục tiêu . Bạn cần một khối đặc biệt cho điều đóBroadcastBlock

Tôi sẽ đề xuất phát sóng tới một khối khác và sử dụng nó đểSubscribe

Nhiệm vụ của BroadcastBlock trong cuộc sống là cho phép tất cả các mục tiêu được liên kết từ khối để có được bản sao của mọi phần tử được xuất bản

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;

Đầu ra

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