【问题标题】:In Azure Durable Functions, how do I appropriately determine the progression of a large number of parallel activities?在 Azure Durable Functions 中,如何正确确定大量并行活动的进程?
【发布时间】:2019-07-15 19:37:16
【问题描述】:

我编写了一个 Durable Function 编排器函数,其主要工作是平均扇出 1,000 个并行活动。由于这些活动的完成是前端用户在技术上等待的事情,因此我希望能够在活动仍在运行时查询进度(在前端显示进度条)。

以下是当前编排器代码的一部分,但我怀疑它是否符合编排器功能的约束 (https://docs.microsoft.com/en-us/azure/azure-functions/durable/durable-functions-checkpointing-and-replay#orchestrator-code-constraints)。

基本上,如果 DF 框架将协调器重放到每个等待,这感觉就像它处理的等待数量不合理:

var replicationTasks = new List<Task<ReplicationOutput>>();
var replicationResults = new List<ReplicationOutput>();
// start up each simulation
for (int i = 0; i < inputs.NumberOfReplications; i++)
{
    var replicationInput = new ReplicationInput();
    var task = context.CallActivityAsync<ReplicationOutput>("SimulationOrchestrator_SimulateReplication", replicationInput);
    replicationTasks.Add(task);
}
// set initial custom status 
var progress = new Progress();
progress.NumberCompleted = 0;
progress.Total = inputs.NumberOfReplications;
progress.TimeStarted = context.CurrentUtcDateTime;
progress.ElapsedTime = context.CurrentUtcDateTime.Subtract(progress.TimeStarted);
context.SetCustomStatus(progress);

// as each task finishes
while (replicationTasks.Any())
{
    Task<ReplicationOutput> nextFinished = await Task.WhenAny(replicationTasks);

    replicationTasks.Remove(nextFinished);
    replicationResults.Add(await nextFinished);

    // update progress object and custom status
    progress.NumberCompleted++;
    progress.ElapsedTime = context.CurrentUtcDateTime.Subtract(progress.TimeStarted);
    context.SetCustomStatus(progress);
}
// aggregate replications together into a single set of results
return new Results(replicationResults);

这在简单的测试条件下不一定会失败,但协调器文档(相当积极地)警告说要保持历史表清晰,避免等待/阻塞等。

是否有记录或“最佳实践”的方法来实现可查询进度的目标?我见过的所有扇出/扇入示例仅使用 await Task.WhenAll(replicationTasks) 仅在 所有 任务完成后继续,我认为这不会允许增量进度检查。

【问题讨论】:

  • 据我所知,这是 DF 中的一个主要差距,并且在某处没有提到最佳实践。对于一个项目,我们通过 Azure Sql Db 中的 Serilog 批量批量插入创建了一个用于大规模扇出 >5k 的状态日志记录。在不同的解决方案中,我们目前正试图通过新的 DurableEntities 实现某种状态。每个活动(按实例)在实体中保存一个状态。因此,您可以实时查询所有实体以了解您的活动。也许这有帮助。

标签: c# asynchronous aggregate monitoring azure-durable-functions


【解决方案1】:

您在这里似乎有两个问题:

大量动作导致重放过多

众所周知,当必须重放大量操作时,持久函数会降低性能。在 Durable Functions 的 .NET 运行时中,DF 在 100,000 次操作后自动中止 (Github) (StackOverflow)。这个 100k 限制是可配置的,但它代表了我发现的“DF 打算处理多少个动作”的唯一准则。

我还没有看到有人在网上讨论减少每个 Durable Function 负责的操作的架构。我有你描述的相同用例(有限并发的扇出),我找到了一些选择:

  1. 使用continueAsNew DF API 调用定期以全新的空重播历史重新启动编排器。该调用带有一个参数,您可以在其中传递您希望 DF 的新实例具有的任何状态(即工作负载的其余部分)。它与函数递归中的概念相同,只是使用了 DF。这是一种直接解决重放性能问题的相对简单的方法,无需在架构中引入新组件。代价是您的处理会定期停止,同时避免产生新的活动以准备重新启动 Orchestrator。

  2. 您可以限制过多的重播批处理任务并使用单独的 DF 层来处理批处理任务。这使得您的顶级 Orchestrator 负责 1/n 的操作,其中 n 是批量大小。我觉得这是一个笨拙的解决方案,它限制了您的顶级 DF 可以监督扇出的粒度,但它解决了重放问题。

  3. 您可以使用扩展会话来延迟 Durable Functions 在调用操作后的一段时间内关闭它们。如果您的 DF 在此期间被激活,它将继续执行而无需重播,就好像它只是一个普通程序在做它的事情一样。您将支付让它一直运行的成本,但如果它有大量的重放历史,它可能会经常醒来,以至于额外的执行成本可以忽略不计。

  4. 如果您不需要对并发和任务调度进行细粒度控制,您可以考虑让 Orchestrator 将任务推送到存储队列,让队列消息触发普通函数。然后,您的 DF 可以在将任务数据移交给队列后立即结束其生命周期。如果它需要处理任务输出,它可以等待恢复处理,直到另一个函数发送一个事件通知它所有任务都已处理完毕。

查询你的DF进度

Durable Functions 有一个明确的机制,可以将其进度传达给相关方。或者,您的 Orchestrator 可以将其进度保存在其他各方可以访问的地方。

您的第一个选择是使用 Durable Function 的 custom status 功能。您可以定期更新 Durable Function 的状态以反映其进度,以及query the DF's progress

另一个选项是持久实体;它们是 DF 存储在 DF 的上下文和生命周期之外持续存在的数据的一种便捷方式,而且至关重要的是,DF 函数应用程序中的客户端函数可以读取 DF 写入的持久实体。围绕它包裹一个 HTTP 触发器,然后bam,你就可以进行进度查询了。

最后,通常通过将进度写入任何主要数据库并从轮询 HTTP 端点中的数据库读回来处理此任务。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-05-20
    • 2020-03-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-01-31
    • 2022-10-18
    相关资源
    最近更新 更多