监控模式是工作流程中的一个循环过程,用于轮询外部系统直到满足某个条件。 例如,它会检查作业状态直到完成,或者观察天气数据直到天空晴朗。 与固定计划计时器触发器不同,监视器在迭代(避免重叠)之间等待,支持动态间隔,并在满足条件或超时过期后自行终止。
本文介绍了如何使用持久化业务流程协调来实现监视器模式。
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进行可配置轮询间隔的作业状态监控。
先决条件
此方案尚没有适用于 Durable Functions 的 PowerShell 示例。
此场景中尚无可用的Durable Functions的Java示例。 请参阅持久任务 SDKs 选项卡。
- .NET 8.0 SDK 或更高版本
- 访问 Azure Durable Task Scheduler 或本地模拟器
- Node.js 22 或更高版本
- 访问 Azure Durable Task Scheduler 或本地模拟器
- Python 3.9 或更高版本
- 访问 Azure Durable Task Scheduler 或本地模拟器
此示例适用于 .NET、JavaScript、Java 和 Python。
- Java 11 或更高版本
- 访问 Azure Durable Task Scheduler 或本地模拟器
监控方案概述
此示例监视某个地点的当前天气状况,如果是晴天,则通过短信通知用户。 可以使用常规的计时器触发函数来检查天气和发送提醒。 但是,此方法的一个问题是 生存期管理。 如果只应发送一条提醒,则在检测到晴天后,监视器需要自我禁用。
此示例监视某个地点的当前天气状况,如果是晴天,则通过短信通知用户。 可以使用常规的计时器触发函数来检查天气和发送提醒。 但是,此方法的一个问题是 生存期管理。 如果只应发送一条提醒,则在检测到晴天后,监视器需要自我禁用。
此示例监视GitHub存储库中的问题个数,并在有超过3个未关闭问题时向用户发出警报。 可以使用定时器触发的常规函数来定期获取打开问题的计数。 但是,此方法的一个问题是 生存期管理。 如果只发送一个警报,监视器需要在检测到 3 个或多个问题后禁用自身。
此方案尚没有适用于 Durable Functions 的 PowerShell 示例。
此场景中尚无可用的Durable Functions的Java示例。 请参阅持久任务 SDKs 选项卡。
监视模式可以结束自身的执行,同时还具有其他优点:
- 监控器按间隔运行,而不是按计划:计时器触发器每小时运行一次;监控器在每次操作之间等待一小时。 监视器的操作不会重叠,除非你另有说明,这对长期执行的任务很重要。
- 监视器可以使用动态时间间隔:可以根据某种条件更改等待时间。
- 监视器可以在满足某种条件时终止,或者由其他进程终止。
- 监视器可以接受参数。 此示例演示如何将相同的监视过程应用于任何请求的位置、电话号码或存储库。
- 监视器可缩放。 由于每个监视器都是编排实例,你可以创建多个监视器,而无需创建新功能或定义更多代码。
- 监视器可轻松集成到更大的工作流。 监控器可以是更复杂的编排函数的一个部分,也可以是 子编排。
此示例监视长时间运行的作业的状态,并在作业完成或超时时返回最终结果。可以使用常规轮询循环来检查作业状态,但此方法在生存期管理和可靠性方面存在限制。
监视模式提供以下优势:
-
持久轮询:编排在进程重启后保持活跃,即使进程失败也能继续监控。
-
可配置间隔:你可以动态调整状态检查之间的等待时间。
-
超时支持:当满足条件或超时过期时,监视器可以终止。
-
状态可见性:客户端可以查询业务流程的自定义状态以查看当前监视进度。
-
可伸缩性:多个监视器可以同时运行,每个监视器都跟踪不同的作业。
Configuration
配置天气 API
C#/JavaScript 示例调用天气 API 来检查当前条件。 需要提供自己的天气 API 密钥并相应地更新示例代码。 示例代码引用了一个 WeatherUndergroundApiKey 应用设置——用你选定的天气提供商的密钥替换这个密钥。
| 应用设置名称 |
数值描述 |
|
WeatherUndergroundApiKey |
您的天气 API 密钥(如有需要,请替换为提供商的密钥名称)。 |
配置天气 API
C#/JavaScript 示例调用天气 API 来检查当前条件。 需要提供自己的天气 API 密钥并相应地更新示例代码。 示例代码引用了一个 WeatherUndergroundApiKey 应用设置——用你选定的天气提供商的密钥替换这个密钥。
| 应用设置名称 |
数值描述 |
|
WeatherUndergroundApiKey |
您的天气 API 密钥(如有需要,请替换为提供商的密钥名称)。 |
Durable Functions的Python示例尚不适用于此方案。
此方案尚没有适用于 Durable Functions 的 PowerShell 示例。
此场景中尚无可用的Durable Functions的Java示例。 请参阅持久任务 SDKs 选项卡。
协调器
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 传递给编排函数。
E3_Monitor函数使用编排函数的标准function.json。
{
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
],
"disabled": false
}
用于实现该函数的代码如下:
const df = require("durable-functions");
const moment = require("moment");
module.exports = df.orchestrator(function* (context) {
const input = context.df.getInput();
context.log(
"Received monitor request. location: " +
(input ? input.location : undefined) +
". phone: " +
(input ? input.phone : undefined) +
"."
);
verifyRequest(input);
const endTime = moment.utc(context.df.currentUtcDateTime).add(6, "h");
context.log(
"Instantiating monitor for " +
input.location.city +
", " +
input.location.state +
". Expires: " +
endTime +
"."
);
while (moment.utc(context.df.currentUtcDateTime).isBefore(endTime)) {
// Check the weather
context.log(
"Checking current weather conditions for " +
input.location.city +
", " +
input.location.state +
" at " +
context.df.currentUtcDateTime +
"."
);
const isClear = yield context.df.callActivity("E3_GetIsClear", input.location);
if (isClear) {
// It's not raining! Or snowing. Or misting. Tell our user to take advantage of it.
context.log(
"Detected clear weather for " +
input.location.city +
", " +
input.location.state +
". Notifying " +
input.phone +
"."
);
yield context.df.callActivity("E3_SendGoodWeatherAlert", input.phone);
break;
} else {
// Wait for the next checkpoint
const nextCheckpoint = moment.utc(context.df.currentUtcDateTime).add(30, "s");
context.log(
"Next check for " +
input.location.city +
", " +
input.location.state +
" at " +
nextCheckpoint.toString()
);
yield context.df.createTimer(nextCheckpoint.toDate()); // accomodate cancellation tokens
}
}
context.log("Monitor expiring.");
});
function verifyRequest(request) {
if (!request) {
throw new Error("An input object is required.");
}
if (!request.location) {
throw new Error("A location input is required.");
}
if (!request.phone) {
throw new Error("A phone number input is required.");
}
}
E3_Monitor函数使用编排函数的标准function.json。
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "context",
"type": "orchestrationTrigger",
"direction": "in"
}
]
}
用于实现该函数的代码如下:
import azure.durable_functions as df
from datetime import timedelta
from typing import Dict
def orchestrator_function(context: df.DurableOrchestrationContext):
monitoring_request: Dict[str, str] = context.get_input()
repo_url: str = monitoring_request["repo"]
phone: str = monitoring_request["phone"]
# Expiration of the repo monitoring
expiry_time = context.current_utc_datetime + timedelta(minutes=5)
while context.current_utc_datetime < expiry_time:
# Count the number of issues in the repo (the GitHub API caps at 30 issues per page)
too_many_issues = yield context.call_activity("E3_TooManyOpenIssues", repo_url)
# If we detect too many issues, we text the provided phone number
if too_many_issues:
# Extract URLs of GitHub issues, and return them
yield context.call_activity("E3_SendAlert", phone)
break
else:
# Reporting the number of statuses found
status = f"The repository does not have too many issues, for now ..."
context.set_custom_status(status)
# Schedule a new "wake up" signal
next_check = context.current_utc_datetime + timedelta(minutes=1)
yield context.create_timer(next_check)
return "Monitor completed!"
main = df.Orchestrator.create(orchestrator_function)
编排器功能会定期检查某个条件,并在满足时发送警报。 它使用持久计时器来控制轮询间隔,并持续运行直到监控器到期。
param($Context)
$input = $Context.Input | ConvertFrom-Json
$expirationTime = (Get-Date).AddHours(6)
$pollingInterval = New-TimeSpan -Seconds $input.pollingIntervalSeconds
while ((Get-Date) -lt $expirationTime) {
# Check current conditions
$isClear = Invoke-DurableActivity -FunctionName 'E3_GetIsClear' -Input $input.location
if ($isClear) {
# Condition met - send alert and exit
Invoke-DurableActivity -FunctionName 'E3_SendGoodWeatherAlert' -Input $input.phone
break
}
# Wait for the next polling interval
Start-DurableTimer -Duration $pollingInterval
}
@FunctionName("E3_Monitor")
public void monitorOrchestrator(
@DurableOrchestrationTrigger(name = "ctx") TaskOrchestrationContext ctx) {
MonitorRequest input = ctx.getInput(MonitorRequest.class);
Instant expirationTime = ctx.getCurrentInstant().plus(Duration.ofHours(6));
int pollingInterval = input.getPollingIntervalSeconds();
while (ctx.getCurrentInstant().isBefore(expirationTime)) {
// Check current conditions
boolean isClear = ctx.callActivity(
"E3_GetIsClear", input.getLocation(), boolean.class).await();
if (isClear) {
// Condition met - send alert and exit
ctx.callActivity("E3_SendGoodWeatherAlert", input.getPhone()).await();
break;
}
// Wait for the next polling interval
Instant nextCheck = ctx.getCurrentInstant().plus(
Duration.ofSeconds(pollingInterval));
ctx.createTimer(nextCheck).await();
}
}
此业务流程协调程序函数执行以下操作:
- 获取 MonitorRequest,其中包含要监视的地点和用于发送短信通知的电话号码(或 Python 示例中的存储库)。
- 确定监视器的过期时间。 为简便起见,本示例使用了硬编码值。
- 调用状态检查活动以确定是否满足条件。
- 如果满足条件,则调用警报功能以发送通知。
- 创建一个持久计时器,以便在下一个轮询间隔恢复编排。 为简便起见,本示例使用了硬编码值。
- 继续运行,直到当前 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
};
}
}
import {
OrchestrationContext,
TOrchestrator,
} from "@microsoft/durabletask-js";
const monitorOrchestrator: TOrchestrator = async function* (
ctx: OrchestrationContext,
input: { jobId: string; pollingIntervalSeconds: number; timeoutSeconds: number }
): any {
const { jobId, pollingIntervalSeconds, timeoutSeconds } = input;
const expirationTime = new Date(
ctx.currentUtcDateTime.getTime() + timeoutSeconds * 1000
);
let checkCount = 0;
while (ctx.currentUtcDateTime < expirationTime) {
// Check current job status
const jobStatus: any = yield ctx.callActivity(checkJobStatus, {
jobId,
checkCount,
});
checkCount = jobStatus.checkCount;
// Make job status available via custom status
ctx.setCustomStatus(jobStatus);
if (jobStatus.status === "Completed") {
return {
jobId,
finalStatus: "Completed",
checksPerformed: checkCount,
};
}
// Wait for next polling interval
yield ctx.createTimer(pollingIntervalSeconds);
}
// Timeout reached
return {
jobId,
finalStatus: "Timeout",
checksPerformed: checkCount,
};
};
import datetime
from durabletask import task
def monitoring_job_orchestrator(ctx: task.OrchestrationContext, job_data: dict) -> dict:
"""
Orchestrator that demonstrates the monitoring pattern.
Periodically checks the status of a job until it completes or times out.
"""
job_id = job_data.get("job_id")
polling_interval = job_data.get("polling_interval_seconds", 5)
timeout = job_data.get("timeout_seconds", 30)
# Record the start time
start_time = ctx.current_utc_datetime
expiration_time = start_time + datetime.timedelta(seconds=timeout)
# Initialize monitoring state
job_status = {
"job_id": job_id,
"status": "Unknown",
"check_count": 0
}
# Loop until the job completes or times out
while True:
# Check current job status
check_input = {"job_id": job_id, "check_count": job_status.get("check_count", 0)}
job_status = yield ctx.call_activity("check_job_status", input=check_input)
# Make the job status available via custom status
ctx.set_custom_status(job_status)
if job_status["status"] == "Completed":
break
# Check if we've hit the timeout
current_time = ctx.current_utc_datetime
if current_time >= expiration_time:
job_status["status"] = "Timeout"
break
# Calculate next check time
next_check_time = current_time + datetime.timedelta(seconds=polling_interval)
if next_check_time > expiration_time:
next_check_time = expiration_time
# Wait until next polling interval
yield ctx.create_timer(next_check_time)
# Return the final status
return {
"job_id": job_id,
"final_status": job_status["status"],
"checks_performed": job_status["check_count"]
}
此示例适用于 .NET、JavaScript、Java 和 Python。
import com.microsoft.durabletask.*;
import com.microsoft.durabletask.azuremanaged.DurableTaskSchedulerWorkerExtensions;
import java.time.Duration;
DurableTaskGrpcWorker worker = DurableTaskSchedulerWorkerExtensions.createWorkerBuilder(connectionString)
.addOrchestration(new TaskOrchestrationFactory() {
@Override
public String getName() { return "MonitoringJobOrchestrator"; }
@Override
public TaskOrchestration create() {
return ctx -> {
JobData jobData = ctx.getInput(JobData.class);
int pollingCount = 0;
// Set initial status
ctx.setCustomStatus(new JobStatus("Starting monitoring..."));
while (true) {
// Update status
ctx.setCustomStatus(new JobStatus(
"Polling job status (attempt " + (++pollingCount) + ")"));
// Wait for polling interval
ctx.createTimer(Duration.ofSeconds(jobData.pollingIntervalSeconds)).await();
// Check if job is complete (simulated after 3 attempts)
if (pollingCount >= 3) {
ctx.setCustomStatus(new JobStatus("Job completed successfully"));
ctx.complete(new JobResult(
"COMPLETED",
"Job completed after " + pollingCount + " attempts"));
break;
}
}
};
}
})
.build();
此编排器执行以下操作:
- 将作业 ID、轮询间隔和超时作为输入参数。
- 记录开始时间并计算过期时间。
- 进入一个检查作业状态的轮询循环。
- 更新自定义状态,以便客户端可以监视进度。
- 如果作业完成,则返回最终结果。
- 如果达到超时值,则返回超时状态。
- 使用
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);
}
function.json 定义如下:
{
"bindings": [
{
"name": "location",
"type": "activityTrigger",
"direction": "in"
}
],
"disabled": false
}
实现如下所示。
const request = require("request");
const clearWeatherConditions = [
"Overcast",
"Clear",
"Partly Cloudy",
"Mostly Cloudy",
"Scattered Clouds",
];
module.exports = function (context, location) {
getCurrentConditions(location)
.then(function (data) {
const isClear = clearWeatherConditions.includes(data.weather);
context.done(null, isClear);
})
.catch(function (err) {
context.log(`E3_GetIsClear encountered an error: ${err}`);
context.done(err);
});
};
function getCurrentConditions(location) {
return new Promise(function (resolve, reject) {
const options = {
url: `https://api.wunderground.com/api/${process.env["WeatherUndergroundApiKey"]}/conditions/q/${location.state}/${location.city}.json`,
method: "GET",
json: true,
};
request(options, function (err, res, body) {
if (err) {
reject(err);
}
if (body.error) {
reject(body.error);
}
if (body.response.error) {
reject(body.response.error);
}
resolve(body.current_observation);
});
});
}
E3_TooManyOpenIssues函数会获得仓库中当前未解决问题的列表,并判断是否“过多”:根据样本,超过3个。
function.json 定义如下:
{
"scriptFile": "__init__.py",
"bindings": [
{
"name": "repoID",
"type": "activityTrigger",
"direction": "in"
}
]
}
实现如下所示。
import requests
import json
def main(repoID: str) -> str:
# We use the GitHub API to count the number of open issues in the repo provided
# Note that the GitHub API only displays at most 30 issues per response, so
# the maximum number this activity will return is 30. That's enough for demo'ing purposes.
[user, repo] = repoID.split("/")
url = f"https://api.github.com/repos/{user}/{repo}/issues?state=open"
res = requests.get(url)
if res.status_code != 200:
error_message = f"Could not find repo {user} under {repo}! API endpoint hit was: {url}"
raise Exception(error_message)
issues = json.loads(res.text)
too_many_issues: bool = len(issues) >= 3
return too_many_issues
E3_GetIsClear函数检查条件是否满足。 在这个例子中,它检查了某个地点的天气状况。
param($location)
# Call weather API to check current conditions
# In a real app, call an external API here
$conditions = Get-WeatherConditions -Location $location
$conditions -eq 'Clear'
E3_GetIsClear函数通过天气API检查某地当前的天气状况。
@FunctionName("E3_GetIsClear")
public boolean getIsClear(
@DurableActivityTrigger(name = "location") Location location) {
// Call weather API to check current conditions
String conditions = getWeatherConditions(location);
return conditions.equals("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 触发函数,您可以通过发送以下 HTTP POST 请求来启动编排:
POST https://{host}/orchestrators/E3_Monitor
Content-Length: 77
Content-Type: application/json
{ "location": { "city": "Redmond", "state": "WA" }, "phone": "+1425XXXXXXX" }
需要一个 GitHub 帐户。 创建可向其提出问题的临时公共存储库。
通过使用示例中包含的 HTTP 触发函数,您可以通过发送以下 HTTP POST 请求来启动编排:
POST https://{host}/orchestrators/E3_Monitor
Content-Length: 77
Content-Type: application/json
{ "repo": "<your GitHub handle>/<a new GitHub repo under your user>", "phone": "+1425XXXXXXX" }
例如,如果你的 GitHub 用户名是 foo,仓库是 bar,请将 "repo" 的值设置为 "foo/bar"。
通过使用示例中包含的 HTTP 触发函数,您可以通过发送以下 HTTP POST 请求来启动编排:
POST https://{host}/api/orchestrators/E3_Monitor
Content-Type: application/json
{ "location": { "city": "Redmond", "state": "WA" }, "phone": "+1425XXXXXXX" }
通过使用示例中包含的 HTTP 触发函数,您可以通过发送以下 HTTP POST 请求来启动编排:
POST https://{host}/api/StartMonitor
Content-Type: application/json
{ "location": { "city": "Redmond", "state": "WA" }, "phone": "+1425XXXXXXX" }
HTTP 触发器函数会计划业务流程:
@FunctionName("StartMonitor")
public HttpResponseMessage startMonitor(
@HttpTrigger(name = "req", methods = {HttpMethod.POST}) HttpRequestMessage<Optional<String>> req,
@DurableClientInput(name = "durableContext") DurableClientContext durableContext,
final ExecutionContext context) {
DurableTaskClient client = durableContext.getClient();
String instanceId = client.scheduleNewOrchestrationInstance("E3_Monitor", req.getBody().get());
context.getLogger().info("Started monitor orchestration with ID = " + instanceId);
return durableContext.createCheckStatusResponse(req, instanceId);
}
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}
若要运行示例,需要:
启动 Durable Task Scheduler 模拟器 (用于本地开发):
docker run -d -p 8080:8080 -p 8082:8082 --name dts-emulator mcr.microsoft.com/dts/dts-emulator:latest
启动工作器 来注册协调器和活动。
运行客户端 来计划监控编排。
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);
}
import {
DurableTaskAzureManagedClientBuilder,
DurableTaskAzureManagedWorkerBuilder,
} from "@microsoft/durabletask-js-azuremanaged";
const client = new DurableTaskAzureManagedClientBuilder()
.connectionString(connectionString)
.build();
const worker = new DurableTaskAzureManagedWorkerBuilder()
.connectionString(connectionString)
.addOrchestrator(monitorOrchestrator)
.addActivity(checkJobStatus)
.build();
await worker.start();
// Schedule the monitoring orchestration
const input = {
jobId: `job-${Date.now()}`,
pollingIntervalSeconds: 5,
timeoutSeconds: 30,
};
const instanceId = await client.scheduleNewOrchestration(
monitorOrchestrator,
input
);
console.log(`Started monitoring orchestration: ${instanceId}`);
// Wait for completion
const result = await client.waitForOrchestrationCompletion(
instanceId,
true,
60
);
console.log(`Final result: ${result?.serializedOutput}`);
await worker.stop();
await client.stop();
from durabletask.azuremanaged.client import DurableTaskSchedulerClient
import time
client = DurableTaskSchedulerClient(
host_address=endpoint,
secure_channel=endpoint != "http://localhost:8080",
taskhub=taskhub,
token_credential=credential
)
# Schedule the monitoring orchestration
job_data = {
"job_id": "job-123",
"polling_interval_seconds": 5,
"timeout_seconds": 30
}
instance_id = client.schedule_new_orchestration(
monitoring_job_orchestrator,
input=job_data
)
print(f"Started monitoring orchestration: {instance_id}")
# Wait for completion
result = client.wait_for_orchestration_completion(instance_id, timeout=60)
print(f"Final result: {result.serialized_output}")
此示例适用于 .NET、JavaScript、Java 和 Python。
import java.time.Duration;
import java.util.UUID;
DurableTaskClient client = DurableTaskSchedulerClientExtensions
.createClientBuilder(connectionString).build();
// Schedule the monitoring orchestration
JobData jobData = new JobData(
"job-" + UUID.randomUUID().toString(),
5, // polling interval seconds
30 // timeout seconds
);
String instanceId = client.scheduleNewOrchestrationInstance(
"MonitoringJobOrchestrator",
new NewOrchestrationInstanceOptions().setInput(jobData));
System.out.println("Started monitoring orchestration: " + instanceId);
// Wait for completion
OrchestrationMetadata result = client.waitForInstanceCompletion(
instanceId, Duration.ofSeconds(60), true);
System.out.println("Final result: " + result.readOutputAs(JobResult.class).status);
后续步骤
本示例演示了如何使用 Durable Functions 通过 durable 定时器和条件逻辑监控外部源的状态。 下一个示例演示如何使用外部事件和 持久计时器 来处理人工交互。
此示例演示了如何使用 Durable Task SDK 实现具有持久计时器和状态跟踪的监视模式。 想了解更多关于其他图案和特征的信息,请参见: