【问题标题】:f# observable fork and side effectf# observable fork 和副作用
【发布时间】: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 创建两个订阅。您应该尝试在 messagesStream observable 的末尾放置一个发布运算符。
  • 谢谢,我发现了一个发布,你知道有发布功能的 f# 库吗?如果您可以添加答案和代码示例,我会将其标记为答案。谢谢

标签: f# system.reactive reactive-programming


【解决方案1】:

添加一些日志记录的最简单方法是使用Observable.iter,如下所示:

let messagesStream = 
  socket.observable
  |> Observable.map (fun () -> socket.read ()) 
  |> Observable.map (fun m -> deserialize m)
  |> Observable.iter (printfn "%A")

messagesStream |> Observable.add (fun m -> handle m)  

【讨论】:

    【解决方案2】:

    听起来您可以创建一个 observable 来处理从网络读取消息并将其反序列化。假设这是使用标准 Rx 运算符完成的,则应该返回一个推送新的反序列化网络消息的 observable。

    您可以有 2 个订阅者订阅该 observable,一个订阅者使用您想要的任何业务逻辑对新消息做出反应,第二个订阅者记录消息。

    这应该消除从网络多次读取的副作用。推送 2 个反序列化消息的副本不会产生副作用。

    【讨论】:

    • 我用一个不起作用的例子编辑了代码。
    • 我不确定我是否理解您的回答。您的建议不是使用 Observable.map 而是订阅网络然后生成新的 observable?能给个代码例子吗?
    • 看到您发布的代码,您似乎尝试了我在回答中的建议。究竟是什么不工作?
    • @Eric - 您描述的解决方案在 C# 中也无法正常工作。 Rx 查询定义了一个管道——它可能涉及分叉——但是当你订阅一个 observable 时,你会创建一个 observable 的实例,一直到源。如果您有两个订阅,那么您就有两个来源。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-06
    • 1970-01-01
    • 2018-03-14
    • 2018-01-17
    • 1970-01-01
    • 2015-06-08
    相关资源
    最近更新 更多