【问题标题】:What is a good way to wrap a MailboxProcessor into and IObservable in F#?在 F# 中将 MailboxProcessor 包装到 IObservable 中的好方法是什么?
【发布时间】:2020-12-21 14:27:29
【问题描述】:

假设我有一个 MailboxProcessor 接收 AsyncReplyChannel 消息并异步执行它们。

有没有像这样从MailboxProcessor 构建IObservable 的简单方法?

let actor = MailboxProcessor.Start (fun inbox ->
  async {
    let mutable x = 0

    while true do
      let! ch : AsyncReplyChannel<int> = inbox.Receive ()

      ch.Reply x

      x <- x + 1
  })

let obs : IObservable<int> = Observable.ofActor actor 

actor.PostAndReply (fun ch -> ch) // Fires obs too

我想围绕热/冷可观察物等做出一些决定。

这里涉及到:Does MailboxProcessor just duplicate IObservable?

【问题讨论】:

  • 为什么要这样做?邮箱处理器和 Observable 代表了两种不同的计算范式。 MailboxProcessor 不复制 IObservable。如果您想使用反应式编程,您可以使用反应式扩展并将处理器发布到主题。 IObservable 本身并不能提供太多
  • MailboxProcessor 是一个非常低级的构造,可用于实现数据流/CSP 管道或代理。虽然它的级别有点太低了,所以使用例如 TPL 数据流块来实现 CSP 管道或使用 Akka.NET 等来实现代理要容易得多。 MailboxProcessor 不提供代理框架提供的任何功能,但缓冲输入和工作线程除外。
  • 此时,您甚至可以询问您是否需要 MailboxProcessor 或者例如 Channel 和 IAsyncEnumerable 是否可以完成相同的工作 - 他们可以。事实上,它们可以提供更好的隔离
  • @PanagiotisKanavos 问题是我想将有序输入发送到进程(MailboxProcessor 适用于此)但进程的输出可能会以不受控制的间隔返回(IObservable)跨度>

标签: f# system.reactive actor


【解决方案1】:

如果你想获得一个现有的MailboxProcessor 并获得一个IObservable,只要邮箱处理器响应消息就会触发,那么我认为没有办法做到这一点 - 没有钩子会让你检测到这一点。

这样做的方法是在MailboxProcessor 上定义某种包装器。例如,您可以定义NotifyingMailboxProcessor&lt;'Msg, 'Evt&gt;,它在每次处理'Msg 类型的消息时触发'Evt 类型的事件。你可以这样开始:

type NotifyingMailboxProcessor<'Msg, 'Evt>(mbox:MailboxProcessor<'T>) = 
  let evt = Event<'E>()
  member x.OnPostAndReply = evt.Publish
  member x.PostAndReply<'R>(f:AsyncReplyChannel<'R> -> 'T, g:'T -> 'R -> 'E) = 
    let mutable msg = Unchecked.defaultof<_>
    let res = mbox.PostAndReply(fun ch -> msg <- f(ch); msg)
    evt.Trigger(g msg res)
    res

type NotifyingMailboxProcessor =
  static member Start<'Msg, 'Evt>(f) = 
    NotifyingMailboxProcessor<'Msg, 'Evt>(MailboxProcessor.Start(f))

这模拟了标准界面,因此您可以使用NotifyingMailboxProcessor.Start 创建一个。一个微妙的问题是,要包装PostAndReply,你需要给它另一个函数,从发送到邮箱的消息和回复组成的一对中构造事件'Evt(这些可以有不同的类型,所以你必须将结果包装到类型 'Evt 中,就像您必须将多种消息类型包装到一个可区分的联合中一样,通常)。

一个受你激励的例子是:

let actor = NotifyingMailboxProcessor.Start<AsyncReplyChannel<int>, int>(fun inbox ->
  async {
    let mutable x = 0
    while true do
      let! (ch : AsyncReplyChannel<int>) = inbox.Receive ()
      ch.Reply x
      x <- x + 1
  })

actor.OnPostAndReply.Add(printfn "Message: %A")
actor.PostAndReply((fun ch -> ch), (fun msg ans -> ans))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-07-30
    • 1970-01-01
    • 2012-06-15
    • 1970-01-01
    • 1970-01-01
    • 2016-05-16
    • 2016-08-18
    • 1970-01-01
    相关资源
    最近更新 更多