モニター パターン

モニターパターンは、外部システムをポーリングし、条件が満たされるまで繰り返し行うワークフローのプロセスです。 例えば、ジョブが完了するまでその状態を確認したり、空が晴れるまで天気データを監視したりします。 固定スケジュール タイマー トリガーとは異なり、モニターはイテレーション間で待機し (重複を回避)、動的な間隔をサポートし、条件が満たされるかタイムアウトが切れると、それ自体を終了できます。

この記事では、耐久性のあるオーケストレーションを用いてモニターパターンを実装する方法を説明します。

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 かローカル エミュレーターにアクセスする

監視シナリオの概要

このサンプルでは、ある場所の現在の気象条件を監視し、空が晴れたときに SMS でユーザーに通知します。 定期的にタイマーでトリガーされる関数を使って、天気を確認し、アラートを送信できます。 ただし、この方法の 1 つの問題は 、有効期間管理です。 送信する必要があるアラートが 1 つだけの場合は、晴天が検出された後でモニターを無効にする必要があります。

特にメリットがあるのは、監視パターンがそれ自体の実行を終了できることです。

  • モニターはスケジュールではなく間隔で実行されます。タイマー トリガーは 1 時間ごとに 実行されます 。モニターはアクションの間に 1 時間 待機します 。 モニターの動作は特に指定されない限り重複しません。これは長期的なタスクでは重要です。
  • モニターの間隔は動的にすることができます。待機時間は、いくつかの条件に基づいて変化することがあります。
  • モニターは、ある条件が満たされたときに終了することも、別のプロセスによって終了することもできます。
  • モニターは、パラメーターを受け取ることができます。 このサンプルでは、要求された場所、電話番号、またはリポジトリに同じ監視プロセスを適用する方法を示します。
  • モニターはスケーラブルです。 各モニターがオーケストレーションインスタンスであるため、新しい関数を作成したりコードを定義したりせずに複数のモニターを作成できます。
  • モニターは、より大規模なワークフローと簡単に統合できます。 モニターには、より複雑なオーケストレーション関数の 1 つのセクションまたは サブオーケストレーションを指定できます。

このサンプルでは、実行時間の長いジョブの状態を監視し、ジョブが完了またはタイムアウトしたときに最終的な結果を返します。通常のポーリング ループを使用してジョブの状態を確認できますが、この方法には 有効期間の管理信頼性に関する制限があります。

監視パターンには、次の利点があります。

  • 永続的なポーリング: オーケストレーションは、プロセスの再起動後も存続し、障害が発生した後も引き続き監視を継続できます。
  • 設定可能な間隔:ステータスチェック間の待ち時間を動的に調整できます。
  • タイムアウトのサポート: モニターは、条件が満たされたとき、またはタイムアウトの有効期限が切れたときに終了できます。
  • 状態の可視性: クライアントは、オーケストレーションのカスタム状態に対してクエリを実行して、現在の監視の進行状況を確認できます。
  • スケーラビリティ: 複数のモニターを同時に実行し、それぞれ異なるジョブを追跡できます。

コンフィギュレーション

Twilio 統合の構成

このサンプルでは、Twilio サービスを使って、携帯電話に SMS メッセージを送信します。 Azure Functionsは既に Twilio バインディングを介して Twilio をサポートしており、このサンプルではその機能を使用しています。

最初に必要なものは Twilio アカウントです。 https://www.twilio.com/try-twilio で無料で作成できます。 アカウントを作成したら、次の 3 つのアプリ設定を関数アプリに追加します。

アプリ設定の名前 値の説明
TwilioAccountSid Twilio アカウントの SID
TwilioAuthToken Twilio アカウントの認証トークン
TwilioPhoneNumber Twilio アカウントに関連付けられている電話番号。 これは、SMS メッセージの送信に使われます。

Weather API の構成

C#/JavaScript サンプルでは、天気 API を呼び出して現在の状態を確認します。 独自の Weather API キーを指定し、それに応じてサンプル コードを更新する必要があります。 サンプルコードはアプリの WeatherUndergroundApiKey 設定を参照しています。このキーを選択した天気プロバイダーのキーに置き換えてください。

アプリ設定の名前 値の説明
WeatherUndergroundApiKey Weather API キー (必要に応じてプロバイダーのキー名に置き換えてください)。

Orchestrator


分離モデル
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を取得します: 監視対象のlocationおよびSMS通知を送信する電話番号(またはPythonの例ではrepo)。
  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を使って現在の気象状況を取得し、空が晴れているかどうかを判断します。


分離モデル
using Microsoft.Azure.Functions.Worker;

namespace VSSample;

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

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

アラート アクティビティ

E3_SendGoodWeatherAlert機能ではTwilioを使って、エンドユーザーに散歩の良い時間だと通知するSMSメッセージを送信します。


分離モデル
using Microsoft.Azure.Functions.Worker;
using Twilio;
using Twilio.Rest.Api.V2010.Account;
using Twilio.Types;

namespace VSSample;

public static partial class Monitor
{
    [Function("E3_SendGoodWeatherAlert")]
    public static async Task SendGoodWeatherAlert(
        [ActivityTrigger] string phoneNumber)
    {
        string accountSid = Environment.GetEnvironmentVariable("TwilioAccountSid")
            ?? throw new InvalidOperationException("TwilioAccountSid is not configured.");
        string authToken = Environment.GetEnvironmentVariable("TwilioAuthToken")
            ?? throw new InvalidOperationException("TwilioAuthToken is not configured.");
        string fromNumber = Environment.GetEnvironmentVariable("TwilioPhoneNumber")
            ?? throw new InvalidOperationException("TwilioPhoneNumber is not configured.");

        TwilioClient.Init(accountSid, authToken);
        await MessageResource.CreateAsync(
            to: new PhoneNumber(phoneNumber),
            from: new PhoneNumber(fromNumber),
            body: "The weather's clear outside! Go take a walk!");
    }
}

隔離されたワーカーサンプルコードを実行するには、 Twilio NuGetパッケージをインストールしてください。


プロセス内モデル
    [FunctionName("E3_SendGoodWeatherAlert")]
    public static void SendGoodWeatherAlert(
        [ActivityTrigger] string phoneNumber,
        ILogger log,
        [TwilioSms(AccountSidSetting = "TwilioAccountSid", AuthTokenSetting = "TwilioAuthToken", From = "%TwilioPhoneNumber%")]
            out CreateMessageOptions message)
    {
        message = new CreateMessageOptions(new PhoneNumber(phoneNumber));
        message.Body = $"The weather's clear outside! Go take a walk!";
    }

internal class WeatherUnderground
{
    private static readonly HttpClient httpClient = new HttpClient();
    private static IReadOnlyDictionary<string, WeatherCondition> weatherMapping = new Dictionary<string, WeatherCondition>()
    {
        { "Clear", WeatherCondition.Clear },
        { "Overcast", WeatherCondition.Clear },
        { "Cloudy", WeatherCondition.Clear },
        { "Clouds", WeatherCondition.Clear },
        { "Drizzle", WeatherCondition.Precipitation },
        { "Hail", WeatherCondition.Precipitation },
        { "Ice", WeatherCondition.Precipitation },
        { "Mist", WeatherCondition.Precipitation },
        { "Precipitation", WeatherCondition.Precipitation },
        { "Rain", WeatherCondition.Precipitation },
        { "Showers", WeatherCondition.Precipitation },
        { "Snow", WeatherCondition.Precipitation },
        { "Spray", WeatherCondition.Precipitation },
        { "Squall", WeatherCondition.Precipitation },
        { "Thunderstorm", WeatherCondition.Precipitation },
    };

    internal static async Task<WeatherCondition> GetCurrentConditionsAsync(Location location)
    {
        var apiKey = Environment.GetEnvironmentVariable("WeatherUndergroundApiKey");
        if (string.IsNullOrEmpty(apiKey))
        {
            throw new InvalidOperationException("The WeatherUndergroundApiKey environment variable was not set.");
        }

        var callString = string.Format("http://api.wunderground.com/api/{0}/conditions/q/{1}/{2}.json", apiKey, location.State, location.City);
        var response = await httpClient.GetAsync(callString);
        var conditions = await response.Content.ReadAsAsync<JObject>();

        JToken currentObservation;
        if (!conditions.TryGetValue("current_observation", out currentObservation))
        {
            JToken error = conditions.SelectToken("response.error");

            if (error != null)
            {
                throw new InvalidOperationException($"API returned an error: {error}.");
            }
            else
            {
                throw new ArgumentException("Could not find weather for this location. Try being more specific.");
            }
        }

        return MapToWeatherCondition((string)(currentObservation as JObject).GetValue("weather"));
    }

    private static WeatherCondition MapToWeatherCondition(string weather)
    {
        foreach (var pair in weatherMapping)
        {
            if (weather.Contains(pair.Key))
            {
                return pair.Value;
            }
        }

        return WeatherCondition.Other;
    }
}

進行中のサンプルコードを実行するには、 Microsoft.Azure.WebJobs.Extensions.Twilio NuGetパッケージをインストールしてください。


このアクティビティは、ジョブの現在の状態を確認します。 実際のアプリケーションでは、このステップは外部APIやサービスを呼び出します。

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

[DurableTask(nameof(CheckJobStatusActivity))]
public class CheckJobStatusActivity : TaskActivity<CheckJobInput, JobStatus>
{
    private readonly ILogger<CheckJobStatusActivity> _logger;

    public CheckJobStatusActivity(ILogger<CheckJobStatusActivity> logger)
    {
        _logger = logger;
    }

    public override Task<JobStatus> RunAsync(TaskActivityContext context, CheckJobInput input)
    {
        _logger.LogInformation("Checking status for job: {JobId} (check #{CheckCount})",
            input.JobId, input.CheckCount + 1);

        // Simulate job status - completes after 3 checks
        var status = input.CheckCount >= 3 ? "Completed" : "Running";

        return Task.FromResult(new JobStatus
        {
            JobId = input.JobId,
            Status = status,
            CheckCount = input.CheckCount + 1,
            LastCheckTime = DateTime.UtcNow
        });
    }
}

// Data classes
public class JobMonitorInput
{
    public string JobId { get; set; }
    public int PollingIntervalSeconds { get; set; } = 5;
    public int TimeoutSeconds { get; set; } = 30;
}

public class CheckJobInput
{
    public string JobId { get; set; }
    public int CheckCount { get; set; }
}

public class JobStatus
{
    public string JobId { get; set; }
    public string Status { get; set; }
    public int CheckCount { get; set; }
    public DateTime LastCheckTime { get; set; }
}

public class JobMonitorResult
{
    public string JobId { get; set; }
    public string FinalStatus { get; set; }
    public int ChecksPerformed { get; set; }
}

モニターのサンプルを実行する

サンプルに含まれる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 timerと条件付き論理を用いて示しています。 次のサンプルでは、外部イベントと 永続的タイマー を使用して人間の相互作用を処理する方法を示します。

このサンプルでは、Durable Task SDK を使用して、永続的タイマーと状態追跡を使用して監視パターンを実装する方法を示しました。 他のパターンや特徴について詳しく知りたい方は、以下をご覧ください: