监视模式

监控模式是工作流程中的一个循环过程,用于轮询外部系统直到满足某个条件。 例如,它会检查作业状态直到完成,或者观察天气数据直到天空晴朗。 与固定计划计时器触发器不同,监视器在迭代(避免重叠)之间等待,支持动态间隔,并在满足条件或超时过期后自行终止。

本文介绍了如何使用持久化业务流程协调来实现监视器模式。

Tip

本文介绍完整的实现。 关于持久编排的应用场景的概念概述,请参见 《什么是持久任务?》

Durable Functions示例包括天气监测方案(C#/JavaScript)和GitHub问题监视方案(Python)。

注释

Azure Functions 的 Node.js 编程模型版本 4 已正式发布。 v4 模型旨在为 JavaScript 和 TypeScript 开发人员提供更灵活、更直观的体验。 有关 v3 和 v4 之间的差异的详细信息,请参阅 迁移指南

在以下代码片段中,JavaScript (PM4) 表示编程模型 v4,即新体验。

Durable Task SDK的示例展示了通过使用.NET、JavaScript、Python和Java进行可配置轮询间隔的作业状态监控。

先决条件

  • .NET 8.0 SDK 或更高版本
  • 访问 Azure Durable Task Scheduler 或本地模拟器

监控方案概述

此示例监视某个地点的当前天气状况,如果是晴天,则通过短信通知用户。 可以使用常规的计时器触发函数来检查天气和发送提醒。 但是,此方法的一个问题是 生存期管理。 如果只应发送一条提醒,则在检测到晴天后,监视器需要自我禁用。

监视模式可以结束自身的执行,同时还具有其他优点:

  • 监控器按间隔运行,而不是按计划:计时器触发器每小时运行一次;监控器在每次操作之间等待一小时。 监视器的操作不会重叠,除非你另有说明,这对长期执行的任务很重要。
  • 监视器可以使用动态时间间隔:可以根据某种条件更改等待时间。
  • 监视器可以在满足某种条件时终止,或者由其他进程终止。
  • 监视器可以接受参数。 此示例演示如何将相同的监视过程应用于任何请求的位置、电话号码或存储库。
  • 监视器可缩放。 由于每个监视器都是编排实例,你可以创建多个监视器,而无需创建新功能或定义更多代码。
  • 监视器可轻松集成到更大的工作流。 监控器可以是更复杂的编排函数的一个部分,也可以是 子编排

此示例监视长时间运行的作业的状态,并在作业完成或超时时返回最终结果。可以使用常规轮询循环来检查作业状态,但此方法在生存期管理和可靠性方面存在限制。

监视模式提供以下优势:

  • 持久轮询:编排在进程重启后保持活跃,即使进程失败也能继续监控。
  • 可配置间隔:你可以动态调整状态检查之间的等待时间。
  • 超时支持:当满足条件或超时过期时,监视器可以终止。
  • 状态可见性:客户端可以查询业务流程的自定义状态以查看当前监视进度。
  • 可伸缩性:多个监视器可以同时运行,每个监视器都跟踪不同的作业。

Configuration

配置天气 API

C#/JavaScript 示例调用天气 API 来检查当前条件。 需要提供自己的天气 API 密钥并相应地更新示例代码。 示例代码引用了一个 WeatherUndergroundApiKey 应用设置——用你选定的天气提供商的密钥替换这个密钥。

应用设置名称 数值描述
WeatherUndergroundApiKey 您的天气 API 密钥(如有需要,请替换为提供商的密钥名称)。

协调器

using Microsoft.Azure.Functions.Worker;
using Microsoft.DurableTask;
using Microsoft.Extensions.Logging;

namespace VSSample;

public static partial class Monitor
{
    [Function("E3_Monitor")]
    public static async Task Run(
        [OrchestrationTrigger] TaskOrchestrationContext context)
    {
        MonitorRequest input = context.GetInput<MonitorRequest>()
            ?? throw new ArgumentNullException(nameof(context), "An input object is required.");
        VerifyRequest(input);

        ILogger logger = context.CreateReplaySafeLogger("E3_Monitor");
        DateTime endTime = context.CurrentUtcDateTime.AddHours(6);
        logger.LogInformation(
            "Instantiating monitor for {Location}. Expires: {EndTime}.",
            input.Location,
            endTime);

        while (context.CurrentUtcDateTime < endTime)
        {
            logger.LogInformation(
                "Checking current weather conditions for {Location} at {CurrentTime}.",
                input.Location,
                context.CurrentUtcDateTime);

            bool isClear = await context.CallActivityAsync<bool>(
                "E3_GetIsClear",
                input.Location);

            if (isClear)
            {
                await context.CallActivityAsync(
                    "E3_SendGoodWeatherAlert",
                    input.Phone);
                break;
            }

            DateTime nextCheckpoint = context.CurrentUtcDateTime.AddMinutes(30);
            await context.CreateTimer(nextCheckpoint, CancellationToken.None);
        }

        logger.LogInformation("Monitor expiring.");
    }

    private static void VerifyRequest(MonitorRequest request)
    {
        ArgumentNullException.ThrowIfNull(request.Location);
        ArgumentException.ThrowIfNullOrEmpty(request.Phone);
    }
}

public sealed class MonitorRequest
{
    public required Location Location { get; init; }

    public required string Phone { get; init; }
}

public sealed class Location
{
    public required string State { get; init; }

    public required string City { get; init; }

    public override string ToString() => $"{City}, {State}";
}
[FunctionName("E3_Monitor")]
public static async Task Run([OrchestrationTrigger] IDurableOrchestrationContext monitorContext, ILogger log)
{
    MonitorRequest input = monitorContext.GetInput<MonitorRequest>();
    if (!monitorContext.IsReplaying) { log.LogInformation($"Received monitor request. Location: {input?.Location}. Phone: {input?.Phone}."); }

    VerifyRequest(input);

    DateTime endTime = monitorContext.CurrentUtcDateTime.AddHours(6);
    if (!monitorContext.IsReplaying) { log.LogInformation($"Instantiating monitor for {input.Location}. Expires: {endTime}."); }

    while (monitorContext.CurrentUtcDateTime < endTime)
    {
        // Check the weather
        if (!monitorContext.IsReplaying) { log.LogInformation($"Checking current weather conditions for {input.Location} at {monitorContext.CurrentUtcDateTime}."); }

        bool isClear = await monitorContext.CallActivityAsync<bool>("E3_GetIsClear", input.Location);

        if (isClear)
        {
            // It's not raining! Or snowing. Or misting. Tell our user to take advantage of it.
            if (!monitorContext.IsReplaying) { log.LogInformation($"Detected clear weather for {input.Location}. Notifying {input.Phone}."); }

            await monitorContext.CallActivityAsync("E3_SendGoodWeatherAlert", input.Phone);
            break;
        }
        else
        {
            // Wait for the next checkpoint
            var nextCheckpoint = monitorContext.CurrentUtcDateTime.AddMinutes(30);
            if (!monitorContext.IsReplaying) { log.LogInformation($"Next check for {input.Location} at {nextCheckpoint}."); }

            await monitorContext.CreateTimer(nextCheckpoint, CancellationToken.None);
        }
    }

    log.LogInformation($"Monitor expiring.");
}

[Deterministic]
private static void VerifyRequest(MonitorRequest request)
{
    if (request == null)
    {
        throw new ArgumentNullException(nameof(request), "An input object is required.");
    }

    if (request.Location == null)
    {
        throw new ArgumentNullException(nameof(request.Location), "A location input is required.");
    }

    if (string.IsNullOrEmpty(request.Phone))
    {
        throw new ArgumentNullException(nameof(request.Phone), "A phone number input is required.");
    }
}

协调器函数需要一个监控位置和一个电话号码,以便在该位置天气转晴时发送消息。 你将这些数据作为强类型对象 MonitorRequest 传递给编排函数。

此业务流程协调程序函数执行以下操作:

  1. 获取 MonitorRequest,其中包含要监视的地点和用于发送短信通知的电话号码(或 Python 示例中的存储库)。
  2. 确定监视器的过期时间。 为简便起见,本示例使用了硬编码值。
  3. 调用状态检查活动以确定是否满足条件。
  4. 如果满足条件,则调用警报功能以发送通知。
  5. 创建一个持久计时器,以便在下一个轮询间隔恢复编排。 为简便起见,本示例使用了硬编码值。
  6. 继续运行,直到当前 UTC 时间超过监视器的过期时间或发送警报。

你可以通过多次调用编排函数来同时运行多个编排器函数实例。 你可以指定监控地点和发送警报的电话号码。 在等待计时器时,编排器功能不会运行,所以你不会因此收费。

业务流程协调程序定期检查作业的状态,并在作业完成或超时时返回。

using Microsoft.DurableTask;
using System;
using System.Threading.Tasks;

[DurableTask(nameof(MonitoringJobOrchestration))]
public class MonitoringJobOrchestration : TaskOrchestrator<JobMonitorInput, JobMonitorResult>
{
    public override async Task<JobMonitorResult> RunAsync(
        TaskOrchestrationContext context, JobMonitorInput input)
    {
        var jobId = input.JobId;
        var pollingInterval = TimeSpan.FromSeconds(input.PollingIntervalSeconds);
        var expirationTime = context.CurrentUtcDateTime.AddSeconds(input.TimeoutSeconds);

        // Initialize monitoring state
        int checkCount = 0;

        while (context.CurrentUtcDateTime < expirationTime)
        {
            // Check current job status
            var jobStatus = await context.CallActivityAsync<JobStatus>(
                nameof(CheckJobStatusActivity),
                new CheckJobInput { JobId = jobId, CheckCount = checkCount });

            checkCount = jobStatus.CheckCount;

            // Make job status available via custom status
            context.SetCustomStatus(jobStatus);

            if (jobStatus.Status == "Completed")
            {
                return new JobMonitorResult
                {
                    JobId = jobId,
                    FinalStatus = "Completed",
                    ChecksPerformed = checkCount
                };
            }

            // Calculate next check time
            var nextCheck = context.CurrentUtcDateTime.Add(pollingInterval);
            if (nextCheck > expirationTime)
            {
                nextCheck = expirationTime;
            }

            // Wait until next polling interval
            await context.CreateTimer(nextCheck, default);
        }

        // Timeout reached
        return new JobMonitorResult
        {
            JobId = jobId,
            FinalStatus = "Timeout",
            ChecksPerformed = checkCount
        };
    }
}

此编排器执行以下操作:

  1. 将作业 ID、轮询间隔和超时作为输入参数。
  2. 记录开始时间并计算过期时间。
  3. 进入一个检查作业状态的轮询循环。
  4. 更新自定义状态,以便客户端可以监视进度。
  5. 如果作业完成,则返回最终结果。
  6. 如果达到超时值,则返回超时状态。
  7. 使用 CreateTimer 在两次轮询之间等待而不消耗资源。

活动

与其他示例一样,帮助器活动函数是使用 activityTrigger 触发器绑定的正则函数。

状态检查活动

E3_GetIsClear函数通过Weather Underground API获取当前天气状况,并判断天空是否晴朗。

        [FunctionName("E3_GetIsClear")]
        public static async Task<bool> GetIsClear([ActivityTrigger] Location location)
        {
            var currentConditions = await WeatherUnderground.GetCurrentConditionsAsync(location);
            return currentConditions.Equals(WeatherCondition.Clear);
        }

运行监视示例

通过使用示例中包含的 HTTP 触发函数,您可以通过发送以下 HTTP POST 请求来启动编排:

POST https://{host}/orchestrators/E3_Monitor
Content-Length: 77
Content-Type: application/json

{ "location": { "city": "Redmond", "state": "WA" }, "phone": "+1425XXXXXXX" }
HTTP/1.1 202 Accepted
Content-Type: application/json; charset=utf-8
Location: https://{host}/runtime/webhooks/durabletask/instances/f6893f25acf64df2ab53a35c09d52635?taskHub=SampleHubVS&connection=Storage&code={SystemKey}
RetryAfter: 10

{"id": "f6893f25acf64df2ab53a35c09d52635", "statusQueryGetUri": "https://{host}/runtime/webhooks/durabletask/instances/f6893f25acf64df2ab53a35c09d52635?taskHub=SampleHubVS&connection=Storage&code={systemKey}", "sendEventPostUri": "https://{host}/runtime/webhooks/durabletask/instances/f6893f25acf64df2ab53a35c09d52635/raiseEvent/{eventName}?taskHub=SampleHubVS&connection=Storage&code={systemKey}", "terminatePostUri": "https://{host}/runtime/webhooks/durabletask/instances/f6893f25acf64df2ab53a35c09d52635/terminate?reason={text}&taskHub=SampleHubVS&connection=Storage&code={systemKey}"}

E3_Monitor实例启动并查询当前条件。 如果满足条件,它将调用活动函数来发送警报;否则,它将设置计时器。 当计时器过期时,编排流程将会恢复。

可以通过查看 Azure Functions 门户中的函数日志来查看编排的活动。

编排在达到超时时或检测到条件被满足时完成。 还可以在另一个函数中使用 terminate API,或调用在前面的 202 响应中引用的 terminatePostUri HTTP POST Webhook。 使用 Webhook 时,将 {text} 替换为提前终止的原因。 HTTP POST URL 大致如下所示:

POST https://{host}/runtime/webhooks/durabletask/instances/f6893f25acf64df2ab53a35c09d52635/terminate?reason=Because&taskHub=SampleHubVS&connection=Storage&code={systemKey}

若要运行示例,需要:

  1. 启动 Durable Task Scheduler 模拟器 (用于本地开发):

    docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
    
  2. 启动工作器 来注册协调器和活动。

  3. 运行客户端 来计划监控编排。

using System;
using System.Threading.Tasks;

var client = DurableTaskClientBuilder.UseDurableTaskScheduler(connectionString).Build();

// Schedule the monitoring orchestration
var input = new JobMonitorInput
{
    JobId = "job-" + Guid.NewGuid().ToString(),
    PollingIntervalSeconds = 5,
    TimeoutSeconds = 30
};

string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(
    nameof(MonitoringJobOrchestration), input);

Console.WriteLine($"Started monitoring orchestration: {instanceId}");

// Wait for completion while checking status
while (true)
{
    var state = await client.GetInstanceMetadataAsync(instanceId, getInputsAndOutputs: true);

    if (state.RuntimeStatus == OrchestrationRuntimeStatus.Completed ||
        state.RuntimeStatus == OrchestrationRuntimeStatus.Failed)
    {
        Console.WriteLine($"Monitoring completed: {state.ReadOutputAs<JobMonitorResult>().FinalStatus}");
        break;
    }

    Console.WriteLine($"Current status: {state.ReadCustomStatusAs<JobStatus>()?.Status}");
    await Task.Delay(2000);
}

后续步骤

本示例演示了如何使用 Durable Functions 通过 durable 定时器和条件逻辑监控外部源的状态。 下一个示例演示如何使用外部事件和 持久计时器 来处理人工交互。

此示例演示了如何使用 Durable Task SDK 实现具有持久计时器和状态跟踪的监视模式。 想了解更多关于其他图案和特征的信息,请参见: