DataflowBlock ITargetSource.AsObservable () चालू ट्रिगर नहीं ()

Aug 28 2020

मैं एक डेटाफ्लोब्लॉक का उपयोग करने की कोशिश कर रहा हूं और मुझे यूनिट परीक्षण के लिए गुजरने वाली वस्तुओं की जासूसी करने की आवश्यकता है।

ऐसा करने के लिए, मैं अपने 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);
    }
}

इसलिए मूल रूप से मैं ट्रांसफॉर्मब्लॉक के लिए अपने पर्यवेक्षक की सदस्यता लेता हूं, और मुझे उम्मीद है कि प्रत्येक मूल्य मेरे पर्यवेक्षक "मूल्यों" सूची में पंजीकृत हो जाएगा।

लेकिन, जब कि IsCompleteयह सच है, और OnError()सफलतापूर्वक पंजीकरण अपवाद पर सेट है, तो OnNext()विधि को कभी भी कॉल नहीं किया जाता है जब तक कि यह पाइप लाइन का आखिरी ब्लॉक नहीं है ... मैं समझ नहीं सकता कि क्यों, क्योंकि "अगलाब्लॉक" इस स्रोत से जुड़ा हुआ है सफलतापूर्वक डेटा प्राप्त करें, जिससे साबित होता है कि कुछ डेटा ब्लॉक से बाहर निकल रहे हैं।

मुझे जो समझ में आता है, वह AsObservableब्लॉक से बाहर निकलने वाले हर मूल्यों की रिपोर्ट करना है और न केवल उन मूल्यों को जो अन्य जुड़े हुए लोगों द्वारा उपभोग नहीं किए गए हैं ...

मैं क्या गलत कर रहा हूं ?

जवाब

2 00110001 Aug 28 2020 at 17:08

आपके द्वारा उन्हें पढ़ने का मौका मिलने से पहले आपके संदेशों का उपभोग किया जा रहा _nextBlockहै।

यदि आप इस लाइन पर टिप्पणी करते हैं तो यह _block.LinkTo(_nextBlock);संभव होगा।

AsObservableएकमात्र उद्देश्य केवल एक ब्लॉक को आरएक्स से खपत करने की अनुमति देना है । यह संदेशों को कई लक्ष्यों तक प्रसारित करने के लिए ब्लॉक के आंतरिक कामकाज को नहीं बदलता है । उसके लिए आपको एक विशेष ब्लॉक की आवश्यकता हैBroadcastBlock

मैं एक और ब्लॉक को प्रसारित करने और उस का उपयोग करने का सुझाव दूंगाSubscribe

ब्रॉडकास्टब्लॉक का जीवन में मिशन है कि ब्लॉक से जुड़े सभी लक्ष्यों को प्रकाशित प्रत्येक तत्व की एक प्रति प्राप्त करने के लिए सक्षम किया जाए

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