如何在Prefect 2.0中获取已调度流程运行的UUID
Prefect调度部署:提前获取流程运行UUID及代码不退出问题解决
问题概述
- 使用Prefect Python API创建带调度的部署后,无法在流程运行前获取其UUID(无调度触发时可正常获取,需该ID用于后续取消等操作)
- 调用
run_deployment后代码始终不退出,推测与函数的异步特性相关
原代码示例
from prefect import flow, task from prefect.deployments import Deployment, run_deployment from datetime import datetime, date, time, timezone # Import the flow: from script import my_flow # Configure the deployment: deployment_name = "my_deployment" # Create the deployment for the flow: deployment = Deployment.build_from_flow( flow = my_flow, name = deployment_name, version = 1, work_queue_name = "my_queue", ) deployment.apply() def main(): # Schedule a flow run based on the deployment: response = run_deployment( name = "my_flow/" + deployment_name, parameters = {my_param}, scheduled_time = dateutil.parser.isoparse(scheduledDate), flow_run_name = "my_run", ) print(response) if __name__ == "__main__": main() exit()
解决方案
1. 提前获取调度的流程运行UUID
当指定scheduled_time触发部署时,run_deployment不会直接返回完整的流程运行对象。可以通过Prefect客户端API,根据部署信息、运行名称和调度时间精准查询目标流程的UUID:
修改后的完整代码:
from prefect import flow, task from prefect.deployments import Deployment, run_deployment from datetime import datetime, timezone, timedelta import dateutil.parser from prefect.client import get_client # Import the flow: from script import my_flow # Configure the deployment: deployment_name = "my_deployment" # Create the deployment for the flow: deployment = Deployment.build_from_flow( flow=my_flow, name=deployment_name, version=1, work_queue_name="my_queue", ) deployment.apply() async def fetch_scheduled_flow_run_id(deployment_name, flow_run_name, scheduled_time): async with get_client() as client: # 按条件查询调度的流程运行 flow_runs = await client.read_flow_runs( flow_name="my_flow", deployment_name=deployment_name, flow_run_name=flow_run_name, scheduled_after=scheduled_time - timedelta(minutes=1), scheduled_before=scheduled_time + timedelta(minutes=1) ) return flow_runs[0].id if flow_runs else None async def main(): # 替换为你的实际调度时间和参数 scheduledDate = "2024-05-20T10:00:00+00:00" my_param = {"your_key": "your_value"} # 异步触发调度部署 await run_deployment( name=f"my_flow/{deployment_name}", parameters=my_param, scheduled_time=dateutil.parser.isoparse(scheduledDate), flow_run_name="my_run", ) # 转换调度时间为UTC时区,确保查询匹配 scheduled_time = dateutil.parser.isoparse(scheduledDate).astimezone(timezone.utc) flow_run_id = await fetch_scheduled_flow_run_id(deployment_name, "my_run", scheduled_time) print(f"已调度的流程运行UUID: {flow_run_id}") if __name__ == "__main__": import asyncio asyncio.run(main())
2. 解决代码不退出问题
原代码中同步调用run_deployment会残留未关闭的异步事件循环,导致代码挂起。改用异步方式调用run_deployment,并通过asyncio.run管理事件循环,能确保任务完成后正常退出。
关于手动设置流程运行ID
Prefect不支持手动指定流程运行UUID,所有ID均由系统自动生成,无法通过API自定义,这一点与你查阅文档的结论一致。
内容的提问来源于stack exchange,提问作者Gauthier
相关产品推荐
相关产品推荐

