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