You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从外部(事件/定时器)终止Azure Durable Activity函数

解决方案:终止Durable Activity函数的两种实现方式

核心问题说明

Azure Durable Functions中,TerminateAsync仅会终止Orchestrator实例,无法直接中断正在运行的Activity;Orchestrator也无法向Activity传递有效的取消令牌。要终止Activity,需通过共享状态机制让Activity主动检查终止信号并跳出循环。


方案一:外部Http触发终止

通过Http接口发送终止信号,Activity定期查询共享状态判断是否终止。

步骤1:定义共享状态实体(Azure Table Storage)

使用Azure Table存储终止标记,需安装Microsoft.Azure.WebJobs.Extensions.Storage包:

public class TerminationStatus
{
    public string PartitionKey => "Termination";
    public string RowKey { get; set; } // 用Orchestrator实例ID作为唯一标识
    public bool ShouldTerminate { get; set; }
}

步骤2:改造Activity函数

让Activity定期检查Table中的终止标记,收到信号后跳出循环:

public class ActivityInput
{
    public string Order { get; set; }
    public string InstanceId { get; set; } // 传递Orchestrator实例ID
}

[FunctionName("DurableFunction_Ship")]
public static async Task<string> Ship(
    [ActivityTrigger] string activityInput,
    [Table("TerminationStatus")] CloudTable terminationTable,
    ILogger log)
{
    var input = JsonConvert.DeserializeObject<ActivityInput>(activityInput);
    log.LogInformation($"Processing order: {input.Order}");

    var query = new TableQuery<TerminationStatus>()
        .Where(TableQuery.GenerateFilterCondition("RowKey", QueryComparisons.Equal, input.InstanceId));

    while (true)
    {
        await Task.Delay(5000); // 用Task.Delay替代Thread.Sleep,避免阻塞线程
        log.LogInformation("Simulating long-running process...");

        // 查询终止状态
        var queryResult = await terminationTable.ExecuteQuerySegmentedAsync(query, null);
        var terminationRecord = queryResult.FirstOrDefault();
        
        if (terminationRecord != null && terminationRecord.ShouldTerminate)
        {
            log.LogInformation($"Terminating activity for instance {input.InstanceId}");
            // 清理终止记录(可选)
            await terminationTable.ExecuteAsync(TableOperation.Delete(terminationRecord));
            return $"Order {input.Order} terminated successfully.";
        }
    }
}

步骤3:改造Orchestrator函数

调用Activity时传递Orchestrator实例ID:

[FunctionName("DurableFunction_Orchestrator")]
public static async Task<List<string>> RunOrchestrator(
    [OrchestrationTrigger] DurableOrchestrationContext context)
{
    var outputs = new List<string>();
    var activityInput = new ActivityInput 
    { 
        Order = "shipped", 
        InstanceId = context.InstanceId 
    };
    
    outputs.Add(await context.CallActivityAsync<string>(
        "DurableFunction_Ship", 
        JsonConvert.SerializeObject(activityInput)));
    
    return outputs;
}

步骤4:实现Http终止触发函数

提供Http接口设置终止标记:

[FunctionName("TerminateInstance")]
public static async Task<IActionResult> Run(
    [DurableClient] IDurableOrchestrationClient client,
    [HttpTrigger(AuthorizationLevel.Anonymous, "get", "post", Route = "terminate/{instanceId}")] HttpRequest req,
    [Table("TerminationStatus")] CloudTable terminationTable,
    string instanceId,
    ILogger log)
{
    var reason = req.Query["reason"] ?? "User-initiated termination";
    log.LogInformation($"Sending termination signal to instance {instanceId}: {reason}");

    // 创建或更新终止记录
    var terminationRecord = new TerminationStatus
    {
        RowKey = instanceId,
        ShouldTerminate = true
    };
    await terminationTable.ExecuteAsync(TableOperation.InsertOrReplace(terminationRecord));

    // 可选:同时终止Orchestrator(若不需要等待Activity结束)
    // await client.TerminateAsync(instanceId, reason);

    return new OkObjectResult($"Termination signal sent successfully to instance {instanceId}");
}

方案二:定时自动终止

由Orchestrator设置定时器,到期后自动触发Activity终止。

改造Orchestrator函数

在调用Activity的同时启动定时器,到期后设置终止标记:

[FunctionName("DurableFunction_Orchestrator")]
public static async Task<List<string>> RunOrchestrator(
    [OrchestrationTrigger] DurableOrchestrationContext context,
    [Table("TerminationStatus")] CloudTable terminationTable,
    ILogger log)
{
    var outputs = new List<string>();
    var activityInput = new ActivityInput 
    { 
        Order = "shipped", 
        InstanceId = context.InstanceId 
    };

    // 设置10分钟后触发的定时器
    var timeoutTask = context.CreateTimer(
        context.CurrentUtcDateTime.AddMinutes(10), 
        CancellationToken.None);
    var activityTask = context.CallActivityAsync<string>(
        "DurableFunction_Ship", 
        JsonConvert.SerializeObject(activityInput));

    // 等待Activity完成或定时器到期
    var completedTask = await Task.WhenAny(activityTask, timeoutTask);

    if (completedTask == timeoutTask)
    {
        // 定时器到期,发送终止信号
        log.LogInformation($"Timeout reached for instance {context.InstanceId}, triggering termination");
        var terminationRecord = new TerminationStatus
        {
            RowKey = context.InstanceId,
            ShouldTerminate = true
        };
        await terminationTable.ExecuteAsync(TableOperation.InsertOrReplace(terminationRecord));
        
        // 等待Activity处理终止并返回结果
        outputs.Add(await activityTask);
        timeoutTask.Cancel();
    }
    else
    {
        // Activity正常完成,取消定时器
        timeoutTask.Cancel();
        outputs.Add(await activityTask);
    }

    return outputs;
}

注意事项

  • Activity中的检查间隔(示例为5秒)可根据业务需求调整,间隔越小终止响应越及时,但会增加存储查询开销。
  • 除了Azure Table Storage,也可使用Redis、Azure Blob等分布式存储作为共享状态载体,只要保证Activity和触发函数能访问即可。
  • 确保函数应用配置了正确的存储连接字符串(AzureWebJobsStorage)。

内容的提问来源于stack exchange,提问作者Dmitry K

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 07:58:00