在F#中混合IObservable和Async <'a>

Mik*_*kov 3 f# asynchronous system.reactive

我有IObservable一个库提供,它从外部服务侦听事件:

let startObservable () : IObservable<'a> = failwith "Given"
Run Code Online (Sandbox Code Playgroud)

对于每个收到的事件,我想执行一个返回的动作Async:

let action (item: 'a) : Async<unit> = failwith "Given"
Run Code Online (Sandbox Code Playgroud)

我正在尝试实现一个处理器

let processor () : Async<unit> =
    startObservable()
    |> Observable.mapAsync action
    |> Async.AwaitObservable
Run Code Online (Sandbox Code Playgroud)

我已经弥补mapAsync并且AwaitObservable:理想情况下它们将由一些图书馆提供,到目前为止我找不到它.

额外要求:

  • 应该按顺序执行操作,以便在处理上一个事件时缓冲后续事件.

  • 如果某个操作引发错误,我希望我的处理器完成.否则,它永远不会完成.

  • Async.Start应该尊重通过的取消令牌.

关于我应该使用的图书馆的任何提示?

Bar*_*ski 5

由于您希望将基于推送的模型(IObservable<>)转换为基于拉取的(Async<>),因此您需要排队以缓冲来自可观察数据的数据.如果队列大小有限 - 哪个tbh.应该是为了使整个管道安全,不要过度内存 - 然后还需要缓冲区溢出的策略.

  1. 一种方法是实现一个MailboxProcessor<>自定义observable,它会将数据发布到它.由于MP是本机F#actor实现,因此它能够使用队列进行有序处理以缓冲峰值.
  2. 另一种选择是使用FSharp.Control.AsyncSeq(特别是AsyncSeq.ofObservableBuffered函数),它将observable转换为基于pull的async可枚举 - 在它下面使用第一点的邮箱处理器:

    startObservable()
    |> AsyncSeq.ofObservableBuffered
    |> AsyncSeq.iterAsync action
    
    Run Code Online (Sandbox Code Playgroud)