【问题标题】:Is returning results from MailboxProcessor via Rx a good idea?通过 Rx 从 MailboxProcessor 返回结果是个好主意吗?
【发布时间】:2016-08-18 22:06:54
【问题描述】:

我对下面的代码示例以及人们的想法有点好奇。 我们的想法是从 NetworkStream (~20 msg/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# system.reactive mailboxprocessor


    【解决方案1】:

    我也做过类似的事情,但使用的是普通的 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 在所有其他方面都与事件非常相似,这可能是使用主题而不是普通事件的一个很好的理由。

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

    • 我注意到您的评论指出对 Log 的第一次调用被忽略了。那是因为您仅在此调用发生后才订阅事件处理程序。我认为您可以在这里使用ReplaySubject variation on the Subject idea - 它会在您订阅它时重播所有事件,因此之前发生的事件不会丢失(但缓存是有成本的)。

    总之,我认为使用Subject 可能是一个好主意——它本质上与使用Event 的模式相同(我认为这是从代理公开通知的标准方式),但它可以让你触发@ 987654334@。由于缓存成本,我可能不会使用ReplaySubject - 您只需要确保在触发任何事件之前订阅即可。

    【讨论】:

      猜你喜欢
      • 2020-09-29
      • 2012-10-16
      • 1970-01-01
      • 2021-08-08
      • 1970-01-01
      • 2010-11-05
      • 1970-01-01
      • 2016-02-24
      相关资源
      最近更新 更多