【发布时间】:2018-02-02 18:12:20
【问题描述】:
我有一个 observable 链,初始 observable 来自网络,每次消息准备好读取时都会被触发。然后下一个处理程序读取消息并反序列化它。现在我有一个 observable 的分支,一个是消息处理程序,另一个是记录消息。
问题是因为我使用 observable 我实际上会尝试阅读消息两次。
我知道使用 Event 而不是 Observable 可以解决问题,但是我会遇到垃圾收集问题,这可能会导致套接字无法被收集。
我想到的一个解决方案是插入某种分隔符,结束一个可观察链并创建一个新的,这样的函数是否已经作为 fsharp 或其他库的一部分存在。
这个问题还有其他解决方案吗?
编辑:
无法正常工作的代码示例
let messagesStream =
socket.observable |>
Observable.map (fun () -> socket.read ()) |>
Observable.map (fun m -> deserialize m)
messagesStream |> Observable.add (fun m -> printf m)
messagesStream |> Observable.add (fun m -> handle m)
【问题讨论】:
-
另一方面:您应该将前向管道放在您输入的函数之前,而不是在您管道的值之后。以这种方式排列它们可以更清楚地了解发生了什么。
-
每次订阅 ab observable 管道都会导致订阅源 observable。您正在为
socket.observable创建两个订阅。您应该尝试在messagesStreamobservable 的末尾放置一个发布运算符。 -
谢谢,我发现了一个发布,你知道有发布功能的 f# 库吗?如果您可以添加答案和代码示例,我会将其标记为答案。谢谢
标签: f# system.reactive reactive-programming