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

如何在Python中每5秒获取Azure Durable Orchestrator的活动状态

问题描述

现有如下Azure Durable Orchestration函数代码:

__init__.py

import logging
import json

import azure.functions as func
import azure.durable_functions as df


def orchestrator_function(context: df.DurableOrchestrationContext):
    input_context = context.get_input()
    requestBody = input_context.get('query')
    parallel_tasks = [  context.call_activity("db", requestBody) , context.call_activity("storage",requestBody)]
    status = { 'status' : "started"}
    context.set_custom_status(status)
    outputs =  context.task_all(parallel_tasks)
   
    #Set Custome Status
    status = { 'status' : "completed"}
    context.set_custom_status(status)
    return [outputs]

main = df.Orchestrator.create(orchestrator_function)

function.json

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "context",
      "type": "orchestrationTrigger",
      "direction": "in"
    }
  ]
}

需要实现:异步调用activity函数直至完成,并且每5秒更新并对外暴露其执行状态。

解决方案

要实现周期性更新状态并跟踪activity任务,需利用Durable Orchestration的定时器和任务状态查询能力,修改 orchestrator 函数逻辑如下:

修改后的__init__.py代码:

import logging
import json

import azure.functions as func
import azure.durable_functions as df


def orchestrator_function(context: df.DurableOrchestrationContext):
    input_context = context.get_input()
    requestBody = input_context.get('query')
    
    # 启动两个异步activity任务
    db_task = context.call_activity("db", requestBody)
    storage_task = context.call_activity("storage", requestBody)
    tasks = [db_task, storage_task]
    task_names = ["db", "storage"]
    
    # 初始状态:标记任务已启动
    status = {
        "status": "running",
        "completed_tasks": [],
        "pending_tasks": task_names.copy()
    }
    context.set_custom_status(status)
    
    # 循环检查任务状态,每5秒更新一次
    while not all(task.is_completed for task in tasks):
        # 等待5秒定时器触发
        yield context.create_timer(context.current_utc_datetime.add_seconds(5))
        
        # 更新当前任务状态
        completed = [name for name, task in zip(task_names, tasks) if task.is_completed]
        pending = [name for name, task in zip(task_names, tasks) if not task.is_completed]
        status = {
            "status": "running",
            "completed_tasks": completed,
            "pending_tasks": pending,
            "completed_count": len(completed),
            "total_count": len(tasks)
        }
        context.set_custom_status(status)
    
    # 所有任务完成,获取结果并更新最终状态
    outputs = [task.result for task in tasks]
    final_status = {
        "status": "completed",
        "completed_tasks": task_names,
        "pending_tasks": [],
        "results": outputs
    }
    context.set_custom_status(final_status)
    
    return outputs

main = df.Orchestrator.create(orchestrator_function)

核心逻辑说明

  1. 异步任务启动:单独启动每个activity任务,不直接阻塞等待完成
  2. 周期性状态更新:用context.create_timer创建5秒间隔定时器,每次触发后检查任务完成情况,更新自定义状态
  3. 状态内容设计:自定义状态包含任务运行阶段、已完成/待完成任务列表、进度统计,外部可通过Durable Functions的状态查询接口获取
  4. 任务完成处理:所有任务完成后收集结果,设置最终完成状态并返回结果

外部获取状态方式

外部客户端可通过Durable Functions的实例状态查询接口获取每5秒更新的自定义状态:

  • 针对HTTP触发场景,调用/runtime/webhooks/durabletask/instances/{instanceId}/status接口(需传入对应实例ID)
  • 响应结果中的customStatus字段即为orchestrator中设置的状态内容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:01:40