DataflowBlock ITargetSource.AsObservable () चालू ट्रिगर नहीं ()
मैं एक डेटाफ्लोब्लॉक का उपयोग करने की कोशिश कर रहा हूं और मुझे यूनिट परीक्षण के लिए गुजरने वाली वस्तुओं की जासूसी करने की आवश्यकता है।
ऐसा करने के लिए, मैं अपने 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ब्लॉक से बाहर निकलने वाले हर मूल्यों की रिपोर्ट करना है और न केवल उन मूल्यों को जो अन्य जुड़े हुए लोगों द्वारा उपभोग नहीं किए गए हैं ...
मैं क्या गलत कर रहा हूं ?
जवाब
आपके द्वारा उन्हें पढ़ने का मौका मिलने से पहले आपके संदेशों का उपभोग किया जा रहा _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