Padrão de exibição

O padrão de monitor é um processo recorrente em um fluxo de trabalho que sonda um sistema externo até que uma condição seja cumprida. Por exemplo, ele verifica o status do trabalho até que ele seja concluído, ou observa dados meteorológicos até que o céu esteja limpo. Ao contrário de um gatilho de temporizador com agendamento fixo, um monitor aguarda entre as iterações (evitando sobreposições), dá suporte a intervalos dinâmicos e pode ser encerrado assim que a condição for atendida ou um tempo limite expirar.

Este artigo explica como implementar o padrão de monitor usando orquestrações duráveis.

Dica

Este artigo mostra a implementação completa. Para uma visão conceitual dos casos de uso de orquestração durável, veja O que é Tarefa Durável?

Os exemplos de Durable Functions incluem um cenário de monitoramento meteorológico (C#/JavaScript) e um cenário de monitoramento de problemas GitHub (Python).

Observação

A versão 4 do modelo de programação Node.js para o Azure Functions está disponível em geral. O modelo v4 foi projetado para fornecer uma experiência mais flexível e intuitiva para desenvolvedores JavaScript e TypeScript. Para obter mais informações sobre as diferenças entre v3 e v4, consulte o guia de migração.

Nos snippets de código a seguir, JavaScript (PM4) denota o modelo de programação v4, a nova experiência.

O exemplo dos SDKs de Tarefas Duráveis demonstra o monitoramento do status do trabalho com intervalos de sondagem configuráveis usando .NET, JavaScript, Python e Java.

Pré-requisitos

  • .NET SDK 8.0 ou posterior
  • Acesso ao agendador de tarefas duráveis Azure ou ao emulador local

Visão geral do cenário de monitoramento

Este exemplo monitora as condições meteorológicas atuais de um local e alerta um usuário por SMS quando as condições climáticas melhoram. Você poderia usar uma funcionalidade programada por temporizador para verificar o clima e enviar alertas. No entanto, um problema com essa abordagem é o gerenciamento de tempo de vida. Se apenas um alerta deve ser enviado, o monitor precisa desabilitar-se após a detecção de clima claro.

O padrão de monitoramento pode encerrar sua própria execução, entre outros benefícios:

  • Os monitores são executados em intervalos, não em agendamentos: um disparador com timer executa a cada hora; um monitor aguarda uma hora entre as ações. As ações de um monitor não se sobrepõem, a menos que você especifique o contrário, o que pode ser importante para tarefas de longa duração.
  • Os monitores podem ter intervalos dinâmicos: o tempo de espera pode mudar com base em alguma condição.
  • Os monitores poderão terminar quando alguma condição for atendida ou terminada por outro processo.
  • Monitores podem receber parâmetros. O exemplo mostra como o mesmo processo de monitoramento pode ser aplicado a qualquer local solicitado, número de telefone ou repositório.
  • Monitores são escalonáveis. Como cada monitor é uma instância de orquestração, você pode criar múltiplos monitores sem precisar criar novas funções ou definir mais código.
  • Os monitores integram-se facilmente em fluxos de trabalho maiores. Um monitor pode ser uma seção de uma função de orquestração mais complexa ou uma sub-orquestração.

Este exemplo monitora o status de uma tarefa de longa duração e retorna o resultado final quando a tarefa é concluída ou atinge o tempo limite. Você poderia usar um loop de polling regular para verificar o status da tarefa, mas essa abordagem tem limitações em relação ao gerenciamento de ciclo de vida e à confiabilidade.

O padrão de monitoramento fornece estes benefícios:

  • Polling durável: A orquestração sobrevive a reinicializações de processos e pode continuar monitorando mesmo após falhas.
  • Intervalos configuráveis: Você pode ajustar dinamicamente o tempo de espera entre as verificações de status.
  • Suporte de tempo limite: o monitor pode terminar quando uma condição é atendida ou um tempo limite expira.
  • Visibilidade do status: os clientes podem consultar o status personalizado da orquestração para ver o progresso atual do monitoramento.
  • Escalabilidade: Vários monitores podem ser executados ao mesmo tempo, cada um monitorando diferentes trabalhos.

Configuração

Configurando a integração com o Twilio

Este exemplo envolve o uso do serviço Twilio para enviar mensagens SMS a um telefone celular. Azure Functions já tem suporte para o Twilio por meio da associação Twilio e o exemplo usa esse recurso.

A primeira coisa de que você precisa é uma conta do Twilio. É possível criar uma gratuitamente em https://www.twilio.com/try-twilio. Quando tiver uma conta, adicione as três seguintes configurações de aplicativo ao seu aplicativo de função.

Nome da configuração do aplicativo Descrição do valor
TwilioAccountSid A SID de sua conta do Twilio
TwilioAuthToken O token de autenticação de sua conta do Twilio
TwilioPhoneNumber O número de telefone associado à sua conta do Twilio. Ele é usado para enviar mensagens SMS.

Configurando uma API de clima

Os exemplos de C#/JavaScript chamam uma API meteorológica para verificar as condições atuais. Você precisa fornecer sua própria chave de API meteorológica e atualizar o código de exemplo adequadamente. O código de exemplo faz referência à configuração de aplicativo WeatherUndergroundApiKey - substitua essa chave pela chave do provedor de clima escolhido.

Nome da configuração do aplicativo Descrição do valor
WeatherUndergroundApiKey Sua chave de API de clima (substitua pelo nome da chave do seu provedor conforme necessário).

Orquestrador


Modelo isolado
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}";
}

Modelo em processo
[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.");
    }
}

A função de orquestrador exige um local para monitorar e um número de telefone para enviar uma mensagem quando o tempo fica claro no local. Você passa esses dados para a função do orquestrador como um objeto MonitorRequest fortemente tipado.

Essa função de orquestrador executa as ações a seguir:

  1. Obtém o MonitorRequest que consiste na localização para monitorar e o número de telefone ao qual envia uma notificação por SMS (ou repo para o exemplo de Python).
  2. Determina o tempo de expiração do monitor. O exemplo usa um valor embutido em código para brevidade.
  3. Chama a atividade de verificação de status para determinar se a condição foi atingida.
  4. Se a condição for atendida, chama a atividade de alerta para enviar uma notificação.
  5. Cria um temporizador durável para retomar a orquestração no próximo intervalo de sondagem. O exemplo usa um valor embutido em código para brevidade.
  6. Continua em execução até que o horário UTC atual passe o tempo de expiração do monitor ou um alerta seja enviado.

Você pode executar múltiplas instâncias de funções orquestradoras simultaneamente chamando a função orquestradora várias vezes. Você pode especificar o local para monitorar e o número de telefone para enviar um alerta. A função Orchestrator não está rodando enquanto espera o timer, então você não é cobrado por ela.

O orquestrador verifica periodicamente o status de um trabalho e retorna quando o trabalho é concluído ou chega ao tempo limite.

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
        };
    }
}

Este orquestrador executa as seguintes ações:

  1. Usa a ID do trabalho, o intervalo de sondagem e o tempo limite como parâmetros de entrada.
  2. Registra a hora de início e calcula o tempo de expiração.
  3. Insere um loop de sondagem que verifica o status do trabalho.
  4. Atualiza o status personalizado para que os clientes possam monitorar o progresso.
  5. Se o trabalho for concluído, retornará o resultado final.
  6. Se o tempo limite for atingido, retornará um status de tempo limite.
  7. Usa CreateTimer para aguardar entre tentativas de sondagem sem consumir recursos.

Atividades

Como com outros exemplo, as funções de atividade auxiliares são funções regulares que usam a associação de gatilho activityTrigger.

Atividade de verificação de status

A função E3_GetIsClear obtém as condições climáticas atuais usando a API Weather Underground e determina se o céu está limpo.


Modelo isolado
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;
    }
}

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

Atividade de alerta

A função E3_SendGoodWeatherAlert usa o Twilio para enviar uma mensagem SMS que notifica o usuário final que é um bom momento para uma caminhada.


Modelo isolado
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!");
    }
}

Observação

Para rodar o código de exemplo de trabalhador isolado, instale o Twilio pacote NuGet.


Modelo em processo
    [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;
    }
}

Observação

Para executar o código de exemplo em processo, instale o Microsoft.Azure.WebJobs.Extensions.Twilio pacote NuGet.


A atividade verifica o status atual do trabalho. Em uma aplicação real, essa etapa chama uma API ou serviço externo.

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; }
}

Executar o exemplo de monitoramento

Usando as funções ativadas por HTTP incluídas no exemplo, você pode iniciar a orquestração enviando a seguinte solicitação 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}"}

A instância E3_Monitor inicia e consulta as condições atuais. Se a condição for atendida, ela chamará uma função de atividade para enviar um alerta; caso contrário, ele define um temporizador. Quando o cronômetro expirar, a orquestração será retomada.

É possível visualizar a atividade da orquestração, observando os logs de função no Portal do Azure Functions.

A orquestração é concluída quando o tempo limite é atingido ou a condição é detectada. Você também pode usar a terminate API dentro de outra função ou invocar o webhook HTTP POST terminatePostUri referenciado na resposta de 202 mencionada anteriormente. Para usar o webhook, substitua {text} pelo motivo do encerramento antecipado. O URL HTTP POST é mais ou menos assim:

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

Para executar o exemplo, você precisa:

  1. Inicie o emulador do Agendador de Tarefas Duráveis (para desenvolvimento local):

    docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
    
  2. Inicia o worker para registrar o orquestrador e as atividades.

  3. Inicie o cliente para agendar uma orquestração de monitoramento.

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);
}

Próximas Etapas 

Este exemplo demonstra como usar Durable Functions para monitorar o status de uma fonte externa usando temporizadores duráveis e lógica condicional. O exemplo a seguir mostra como usar eventos externos e temporizadores duráveis para lidar com a interação humana.

Este exemplo demonstrou como usar os SDKs de Tarefa Durável para implementar o padrão de monitoramento com temporizadores duráveis e acompanhamento de status. Para saber mais sobre outros padrões e características, veja: