【问题标题】:What is the recommended F# pattern for this situation?对于这种情况,推荐的 F# 模式是什么?
【发布时间】:2011-08-18 06:37:31
【问题描述】:

我的情况类似于以下:

let mutable stopped = false

let runAsync() = async {
    while not stopped do
        let! item = fetchItemToProcessAsync
        match item with
        | Some job -> job |> runJobAsync |> Async.Start
        | None -> do! Async.Sleep(1000)
}

let run() = Async.Start runAsync
let stop() =
    stopped <- true

现在,当调用 stop 方法时,我必须停止从数据库中读取更多项目,并等待当前启动的项目完成,然后再从该函数返回。

实现这一目标的最佳方法是什么?我正在考虑使用一个计数器,(使用互锁的 API)并在计数器达到 0 时从 stop 方法返回。

如果有其他方法可以做到这一点,我将不胜感激。我有一种感觉,我可以在这里使用代理,但我不确定是否有任何可用的方法可以使用代理完成此操作,或者我是否仍需要编写自定义逻辑来确定作业已完成执行。

【问题讨论】:

    标签: asynchronous f# agent


    【解决方案1】:

    看看actor-based patterns and MailboxProcessor

    基本上你可以把它想象成一个异步队列。如果您使用运行列表(以 Async.StartChildAsync.StartAsTask 开头)作为 MailboxProcessor 内循环的参数,您可以通过等待或 CancellationToken 优雅地处理关闭)

    这是我整理的一个快速示例:

    
    type Commands = 
        | RunJob of Async
        | JobDone of int
        | Quit of AsyncReplyChannel
    
    type JobRunner() =
        let processor =
            MailboxProcessor.Start (fun inbox ->
                let rec loop (nextId, jobs) = async {
                    let! cmd = inbox.Receive()
                    match cmd with
                    | Quit cb ->
                        if not (Map.isEmpty jobs) 
                        then async {
                                do! Async.Sleep 100
                                inbox.Post (Quit cb)}
                            |> Async.Start
                            return! loop (nextId, jobs)
                        else 
                            cb.Reply()
                            return ()
                    | JobDone id ->
                        return! loop (nextId, jobs |> Map.remove id)
                    | RunJob job ->
                        let runJob i = async {
                            do! job
                            inbox.Post (JobDone i)
                        }
                        let! child = Async.StartChild (runJob nextId)
                        return! loop (nextId+1, jobs |> Map.add nextId child)
                }
                loop (0, Map.empty))
        member jr.PostJob(job) = processor.Post (RunJob job)
        member jr.Quit() = processor.PostAndReply(fun cb -> Quit cb)
    
    let postWaitJob (jobRunner : JobRunner) time =
        let job = async {
            do! Async.Sleep time
            printfn "sleept for %d ms" time }
        jobRunner.PostJob job
    
    let testRun() =
        let jr = new JobRunner()
        printfn "starting jobs..."
        [10..-1..1] |> List.iter (fun i -> postWaitJob jr (i*1000))
        printfn "sending quit"
        jr.Quit()
        printfn "done!"
    

    嗯...这里的编辑器有一些问题:当我使用管道返回运算符时,它只会杀死很多代码... grrr

    简短说明:如您所见,我总是为内部循环提供下一个空闲作业 ID 和 Id->AsyncChild 作业的映射。 (您当然可以实施其他/更好的解决方案 - 在此示例中不需要地图,但您可以使用命令“取消 JobNr”或任何其他方式进行扩展) Job done 消息仅在内部用于从该地图中删除作业 退出只是检查映射是否为空 - 如果不需要额外的工作并且邮箱处理器退出(return ()) - 如果它不为空,则启动一个新的 Async-Child,它只等待 100 毫秒,然后重新发送 Quit-Message RunJob 相当简单 - 它只是将给定的作业与 JobDone 的帖子链接到 MessabeboxProcessor 中,然后使用更新的值递归调用循环(nextId 是一个向上,新的 Job 映射到旧的 nextId)

    【讨论】:

    • 感谢您的回复 - 我也在走同样的路。但是,在查看您的示例时,我有两个问题:1)当我们使用Async.StartChild : Async&lt;'u&gt; 开始异步时,我们是否不需要再次调用来检索'u from Async&lt;'u&gt; (e.g. let! x = child or do! child)? 2) 鉴于这些工作可能在不同的时间完成,我们不需要为 Map 函数使用某种同步吗?
    • 嗨 - 是的,如果你想要结果,你会得到另一个让!为孩子 - 但因为我不关心答案(没有)你不必。第二:不,因为我从不访问同一个地图,但总是生成新的地图(地图是不可变的)我不必考虑太多的并发性——除了当前循环当前继续/运行的线程将访问它(这就是为什么我将 JobDone 消息发布到处理器而不是更改全局地图/字典或其他任何内容)
    • 我建议您在其中放置一些“printf”并使用代码 - 在这些情况下总是可以帮助我(如果您运行 testRun 函数(您可以将其全部粘贴到 F# 交互式中)您例如,可以看到在任何任务完成之前调用了退出,但该函数仅在完成所有 10 个作业时才返回)
    • 谢谢,看来该解决方案对我有用!再次感谢!
    【解决方案2】:

    在 fssnip.net 上查看此 snippet。这是您可以使用的通用作业处理器。

    【讨论】:

    • 与我的解决方案几乎相同 - 但我喜欢分成“退出”阶段。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-31
    • 1970-01-01
    • 1970-01-01
    • 2013-04-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多