通过Rx从MailboxProcessor返回结果是一个好主意吗?



我对下面的代码示例和人们的想法有点好奇。这个想法是从一个NetworkStream (~ 20msg/s)中读取,而不是在main中工作,将事情传递给MainboxProcessor来处理,并在完成后将事情返回绑定。

通常的方法是使用PostAndReply,但我想绑定到ListView或c#中的其他控件。无论如何,必须对LastN项和过滤执行魔法。另外,Rx有一些错误处理。

下面的例子从2观察数字。10并返回"hello X"。在8时,它像EOF一样停止。使它成为ToEnumerable,因为其他线程在此之前完成,但它也适用于订阅。

什么让我烦恼:

  1. 递归传递Subject(obj)。我不觉得有3-4个这样的人有什么问题。好主意吗?
  2. 受试者寿命。

open System
open System.Threading
open System.Reactive.Subjects
open System.Reactive.Linq  // NuGet, take System.Reactive.Core also.
open System.Reactive.Concurrency
type SerializedLogger() = 
    let _letters = new Subject<string>()
    // create the mailbox processor
    let agent = MailboxProcessor.Start(fun inbox -> 
        // the message processing function
        let rec messageLoop (letters:Subject<string>) = async{
            // read a message
            let! msg = inbox.Receive()
            printfn "mailbox: %d in Thread: %d" msg Thread.CurrentThread.ManagedThreadId
            do! Async.Sleep 100
            // write it to the log    
            match msg with
            | 8 -> letters.OnCompleted() // like EOF.
            | x -> letters.OnNext(sprintf "hello %d" x)
            // loop to top
            return! messageLoop letters
            }
        // start the loop
        messageLoop _letters
        )
    // public interface
    member this.Log msg = agent.Post msg
    member this.Getletters() = _letters.AsObservable()
/// Print line with prefix 1.
let myPrint1 x = printfn "onNext - %s,  Thread: %d" x  Thread.CurrentThread.ManagedThreadId
// Actions
let onNext = new Action<string>(myPrint1)
let onCompleted = new Action(fun _ -> printfn "Complete")
[<EntryPoint>]
let main argv = 
    async{
    printfn "Main is on: %d" Thread.CurrentThread.ManagedThreadId
    // test
    let logger = SerializedLogger()
    logger.Log 1 // ignored?
    let xObs = logger
                .Getletters() //.Where( fun x -> x <> "hello 5")
                .SubscribeOn(Scheduler.CurrentThread)
                .ObserveOn(Scheduler.CurrentThread)
                .ToEnumerable() // this
                //.Subscribe(onNext, onCompleted) // or with Dispose()
    [2..10] |> Seq.iter (logger.Log) 
    xObs |> Seq.iter myPrint1
    while true 
        do 
        printfn "waiting"
        System.Threading.Thread.Sleep(1000)
    return 0
    } |> Async.RunSynchronously // return an integer exit code

我也做过类似的事情,但使用的是普通的f# Event类型而不是Subject。它基本上允许您创建IObservable并触发其订阅-就像您使用更复杂的Subject一样。基于事件的版本将是:

type SerializedLogger() = 
   let letterProduced = new Event<string>()
   let lettersEnded = new Event<unit>()
   let agent = MailboxProcessor.Start(fun inbox -> 
     let rec messageLoop (letters:Subject<string>) = async {
       // Some code omitted
       match msg with
       | 8 -> lettersEnded.Trigger()
       | x -> letterProduced.Trigger(sprintf "hello %d" x)
       // ...
member this.Log msg = agent.Post msg
member this.LetterProduced = letterProduced.Publish
member this.LettersEnded = lettersEnded.Publish

重要的区别是:

  • Event不能触发OnCompleted,所以我暴露了两个独立的事件。这真是太不幸了!考虑到Subject在所有其他方面都与事件非常相似,这可能是使用subject而不是普通事件的一个很好的理由。

  • 使用Event的好处是它是标准的f#类型,因此您不需要代理中的任何外部依赖项。

  • 我注意到你的评论,注意到第一次调用Log被忽略了。这是因为只有在调用发生后才订阅事件处理程序。我认为你可以在这里使用ReplaySubject的变化主题的想法-它重播所有的事件,当你订阅它,所以一个发生在早些时候不会丢失(但有一个成本缓存)。

总之,我认为使用Subject可能是一个好主意-它本质上与使用Event的模式相同(我认为这是暴露代理通知的相当标准的方式),但它允许您触发OnCompleted。我可能不会使用ReplaySubject,因为缓存成本—您只需要确保在触发任何事件之前订阅。

相关内容

  • 没有找到相关文章

最新更新