如何在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)
核心逻辑说明
- 异步任务启动:单独启动每个activity任务,不直接阻塞等待完成
- 周期性状态更新:用
context.create_timer创建5秒间隔定时器,每次触发后检查任务完成情况,更新自定义状态 - 状态内容设计:自定义状态包含任务运行阶段、已完成/待完成任务列表、进度统计,外部可通过Durable Functions的状态查询接口获取
- 任务完成处理:所有任务完成后收集结果,设置最终完成状态并返回结果
外部获取状态方式
外部客户端可通过Durable Functions的实例状态查询接口获取每5秒更新的自定义状态:
- 针对HTTP触发场景,调用
/runtime/webhooks/durabletask/instances/{instanceId}/status接口(需传入对应实例ID) - 响应结果中的
customStatus字段即为orchestrator中设置的状态内容
内容的提问来源于stack exchange,提问作者microset
相关产品推荐
相关产品推荐

