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应该尊重通过的取消令牌.
关于我应该使用的图书馆的任何提示?
由于您希望将基于推送的模型(IObservable<>)转换为基于拉取的(Async<>),因此您需要排队以缓冲来自可观察数据的数据.如果队列大小有限 - 哪个tbh.应该是为了使整个管道安全,不要过度内存 - 然后还需要缓冲区溢出的策略.
MailboxProcessor<>自定义observable,它会将数据发布到它.由于MP是本机F#actor实现,因此它能够使用队列进行有序处理以缓冲峰值.另一种选择是使用FSharp.Control.AsyncSeq(特别是AsyncSeq.ofObservableBuffered函数),它将observable转换为基于pull的async可枚举 - 在它下面使用第一点的邮箱处理器:
startObservable()
|> AsyncSeq.ofObservableBuffered
|> AsyncSeq.iterAsync action
Run Code Online (Sandbox Code Playgroud)| 归档时间: |
|
| 查看次数: |
239 次 |
| 最近记录: |