动态构建工作流)
Conductor 如何用 Python SDK 以代码方式workflow as code动态构建工作流【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductorConductor 支持 code-first 的工作流构建方式用 Python SDK 以代码定义工作流替代手工编写 JSON。通过运算符链式串联任务并可以加入条件分支Switch、并行Fork/Join、循环Do/While甚至在工作流启动时动态生成任务图。适用前提一个可访问的 Conductor 服务器以及已安装的 Python SDK。准备环境安装 SDK 并配置服务器地址。文档给出的标准方式是设置CONDUCTOR_SERVER_URL环境变量Configuration()会从环境中读取它认证相关变量为CONDUCTOR_AUTH_*pip install conductor-python export CONDUCTOR_SERVER_URLhttp://localhost:8080/apiPython 中获取执行器的标准连接代码from conductor.client.configuration.configuration import Configuration from conductor.client.orkes_clients import OrkesClients config Configuration() # reads CONDUCTOR_SERVER_URL from env clients OrkesClients(configurationconfig) executor clients.get_workflow_executor()一个容易遗漏的前提工作流中的自定义任务SIMPLE 任务需要 worker 轮询执行。SDK 提供TaskHandler它会自动发现所有worker_task装饰的函数并为每个 worker 启动一个子进程from conductor.client.automator.task_handler import TaskHandler with TaskHandler(configurationconfig, scan_for_annotated_workersTrue) as task_handler: task_handler.start_processes()如果 worker 没有在跑工作流会启动但任务一直处于等待状态同步执行会一直阻塞。用 运算符构建顺序工作流用worker_task装饰的普通 Python 函数就是可复用的任务积木task_definition_name是任务的注册名。文档中的示例订单履约流程如下其中的返回值如99.99、txn_abc123是文档示例值实际业务请替换函数体from conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.worker.worker_task import worker_task worker_task(task_definition_namefetch_order) def fetch_order(order_id: str) - dict: return {order_id: order_id, amount: 99.99, item: Widget} worker_task(task_definition_nameprocess_payment) def process_payment(order_id: str, amount: float) - dict: return {transaction_id: txn_abc123, status: charged} worker_task(task_definition_nameship_order) def ship_order(order_id: str, transaction_id: str) - dict: return {tracking: TRACK-456, carrier: FedEx} workflow ConductorWorkflow(nameorder_fulfillment, version1, executorexecutor) fetch fetch_order(task_ref_namefetch, order_idworkflow.input(order_id)) pay process_payment( task_ref_namepay, order_idworkflow.input(order_id), amountfetch.output(amount), ) ship ship_order( task_ref_nameship, order_idworkflow.input(order_id), transaction_idpay.output(transaction_id), ) workflow fetch pay ship workflow.output_parameters({ tracking: ship.output(tracking), transaction_id: pay.output(transaction_id), }) workflow.register(overwriteTrue)几个关键用法workflow.input(order_id)引用工作流启动时的输入字段作为任务输入参数fetch.output(amount)引用上游任务输出中的某个字段作为下游任务的输入任务间的数据传递就这样串起来task_ref_name是任务在流程内的引用名register(overwriteTrue)把工作流定义注册到服务器overwriteTrue表示已存在同名版本时覆盖。把 worker 启动、定义构建和执行放在同一个应用里quickstart.py的结构来自 Python SDK 文档的完整示例from conductor.client.automator.task_handler import TaskHandler from conductor.client.configuration.configuration import Configuration from conductor.client.orkes_clients import OrkesClients from conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.worker.worker_task import worker_task worker_task(task_definition_namegreet, register_task_defTrue) def greet(name: str) - str: return fHello {name} def main(): config Configuration() clients OrkesClients(configurationconfig) executor clients.get_workflow_executor() workflow ConductorWorkflow(namegreetings, version1, executorexecutor) greet_task greet(task_ref_namegreet_ref, nameworkflow.input(name)) workflow greet_task workflow.output_parameters({result: greet_task.output(result)}) workflow.register(overwriteTrue) # 启动 worker 子进程开始轮询 with TaskHandler(configurationconfig, scan_for_annotated_workersTrue) as task_handler: task_handler.start_processes() run executor.execute(namegreetings, version1, workflow_input{name: Conductor}) print(fresult: {run.output[result]}) print(fexecution: {config.ui_host}/execution/{run.workflow_id}) if __name__ __main__: main()注意register_task_defTrue的用途它让 SDK 在本地开发时顺带注册任务定义文档同时提示生产环境中应单独管理任务定义不依赖这个开关。同步执行并验证结果executor.execute会阻塞直到工作流完成返回对象的status和output就是文档给出的验证方式run executor.execute( nameorder_fulfillment, version1, workflow_input{order_id: ORD-789}, ) print(fStatus: {run.status}) print(fOutput: {run.output}) print(fView: {config.ui_host}/execution/{run.workflow_id})判断是否成功的两条路径打印的run.output中应包含output_parameters声明的键上面示例即tracking、transaction_id值对应各 worker 的返回值打开config.ui_host拼接出的执行页面在 Conductor UI 中查看该次执行的每个任务状态这也是文档推荐用于进一步检查的方式。运行时动态生成工作流定义workflow as code 最强的用法是运行时动态工作流不预先注册定义而是用代码在启动时拼装workflow_def随StartWorkflowRequest一起提交适合步骤事先无法确定的场景文档点名了 AI agent 动态生成执行计划这类用途from conductor.client.configuration.configuration import Configuration from conductor.client.orkes_clients import OrkesClients from conductor.client.http.models import StartWorkflowRequest config Configuration() clients OrkesClients(configurationconfig) executor clients.get_workflow_executor() # 步骤列表在运行时确定 steps [validate, enrich, store] tasks [] for i, step in enumerate(steps): tasks.append({ name: step, taskReferenceName: f{step}_{i}, type: SIMPLE, inputParameters: { data: ${workflow.input.data} if i 0 else f${{{steps[i-1]}_{i-1}.output.result}}, }, }) # 内联定义直接启动 —— 无需预先注册 request StartWorkflowRequest( namedynamic_pipeline, workflow_def{ name: dynamic_pipeline, version: 1, tasks: tasks, outputParameters: { result: f${{{steps[-1]}_{len(steps)-1}.output.result}}, }, }, input{data: {key: value}}, ) workflow_id executor.start_workflow(request) print(fStarted dynamic workflow: {workflow_id})两点说明代码里的${workflow.input.data}、${validate_0.output.result}等是 Conductor 的输入表达式语法写在任务inputParameters里由服务器在执行时解析不是 Python 变量或需要人工替换的模板Python f-string 负责在运行时把它们拼成正确的字符串这条路径走start_workflow异步启动立即返回workflow_id。验证方式是拿这个 ID 去 UI 检查执行状态validate、enrich、store对应的 worker 必须已在运行任务才会被消费。条件分支、并行与循环以下三种结构用于替代手工 JSON 时的等价能力示例中的 worker 函数classify_ticket、page_oncall、check_credit等文档未给出函数体需要你用worker_task自行定义后再套用。Switch 条件分支——按任务输出路由每个 case 是一条独立的任务链from conductor.client.workflow.task.switch_task import SwitchTask switch SwitchTask(task_ref_namepriority_router, case_expressionclassify.output(priority)) switch.switch_case(critical, [ page_oncall(task_ref_namepage, ticket_idworkflow.input(ticket_id)), escalate(task_ref_nameescalate, ticket_idworkflow.input(ticket_id)), ]) switch.switch_case(high, [ assign_senior(task_ref_nameassign, ticket_idworkflow.input(ticket_id)), ]) switch.default_case([ add_to_backlog(task_ref_namebacklog, ticket_idworkflow.input(ticket_id)), ]) workflow classify switchFork/Join 并行——ForkTask的forked_tasks是分支列表JoinTask等待所有分支完成之后用.output()合并各分支结果from conductor.client.workflow.task.fork_task import ForkTask from conductor.client.workflow.task.join_task import JoinTask fork ForkTask( task_ref_nameparallel_checks, forked_tasks[[credit_check], [fraud_check], [kyc_check]], ) join JoinTask(task_ref_namewait_all, join_on[credit, fraud, kyc]) workflow fork join decide workflow.output_parameters({decision: decide.output(result)})Do/While 循环——termination_condition是终止表达式配合max_iterations限制迭代次数文档指出它适合轮询、重试和迭代式 AI agent 循环from conductor.client.workflow.task.do_while_task import DoWhileTask loop DoWhileTask( task_ref_nameagent_loop, termination_conditionif ($.act[output][done] true) { false; } else { true; }, tasks[think, act], ) loop.input_parameters.update({max_iterations: 10}) workflow loop summarize限制与边界上面各示例都假设已存在可用的executor见准备环境一节示例中未展示的 worker 函数需要读者按worker_task模式补齐文档不保证其函数体内联workflow_def的工作流不经过register只对本次执行生效不会出现在已注册的工作流定义列表中register(overwriteTrue)会覆盖同名版本的既有定义注意不要用它覆盖线上正在使用的版本完整生命周期操作start/pause/resume/terminate/retry/restart/rerun/signal/search与所有任务类型的更多 Python 示例见 Python SDK 文档各模式的原始示例见 Dynamic workflows in code。【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考