
Conductor 动态工作流实战用 Python 以代码方式构建运行时自适应的任务编排【免费下载链接】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/conductoroutput文章output文章Conductor 动态工作流实战用 Python 以代码方式构建运行时自适应的任务编排Conductor 支持代码优先code-first的工作流定义方式开发者可以在 Python 中直接编排任务图用操作符串联任务、用 Switch 表达条件分支、用 Fork/Join 实现并行、用 DoWhile 实现循环甚至可以在运行时现场生成整个工作流定义并即刻执行。本文以 Conductor 官方 Cookbook 中 dynamic-workflows.md 为核心骨架结合仓库源码深入讲解这套 Workflow as code 的完整玩法。读完本文你将掌握顺序/分支/并行/循环/子工作流等全部代码化编排手段以及为 AI Agent、数据管道按运行时计划动态生成任务图的实战方案。Workflow as code为什么用代码写工作流传统做法是用手写 JSON 定义工作流。而 Conductor 的 Python SDK 提供了另一种思路——用原生 Python 代码构建工作流即 Workflow as code用worker_task装饰器把普通 Python 函数变成可复用的任务构建块用操作符把任务串联成执行链代码结构即执行顺序条件逻辑、循环、并行分支全部用 Python 对象表达天然支持版本管理与代码评审最关键的场景是动态工作流当任务图需要在运行时才能确定时例如取决于 API 响应、用户输入或 LLM 生成的执行计划代码化构建远比维护 JSON 字符串灵活。从 docs/devguide/concepts/workflows.md 的定义看Conductor 把 动态工作流 列为区别于传统编排方案的五大特性之一Workflows can be created and modified at runtime as code-first or JSON definitions, enabling use cases where the task graph is not known ahead of time.本文讨论的正是这条路径。最小顺序工作流先用一个订单履约order_fulfillment示例理解最基本的代码化写法三个worker_task函数分别表示拉取订单、处理支付、发货然后用workflow fetch pay ship串联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)要点拆解task_ref_name是任务的引用名后续所有xxx.output(...)都依赖它取前序任务输出workflow.input(order_id)引用工作流入参fetch.output(amount)引用任务输出这正是工作流定义中${...}表达式在 Python 侧的等价物参见 docs/devguide/concepts/tasks.md 中关于任务配置与${...}数据传递的说明workflow.output_parameters({...})声明工作流输出workflow.register(overwriteTrue)把定义注册到 Conductor 服务端overwriteTrue允许覆盖同名同版本定义。条件分支Switch 任务当执行路径取决于任务输出或工作流入参时使用SwitchTask。每个分支case对应一条独立的任务链未命中任何 case 时走default_casefrom conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.workflow.task.switch_task import SwitchTask workflow ConductorWorkflow(nameroute_by_priority, version1, executorexecutor) classify classify_ticket( task_ref_nameclassify, descriptionworkflow.input(description), ) switch SwitchTask(task_ref_namepriority_router, case_expressionclassify.output(priority)) # Each case is a list of tasks to execute 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 switch workflow.register(overwriteTrue)服务端如何执行 Switch从源码 SwitchTaskMapper.java 可以看到完整决策过程先根据evaluatorType从注册表中取出求值器Evaluator对expression求值若没有注册对应求值器工作流会直接以TerminateWorkflowException终止求值结果写入 Switch 任务输出case/evaluationResult/selectedCase以求值结果作为 key 查decisionCases取出要执行的分支任务列表关键细节只有当decisionCases.get(evalResult)返回null即没有任何 case 匹配时才回退到defaultCase若某 case 显式定义为空列表则视为该分支故意不执行任何任务不会误触发 default。这与 Python SDK 中default_case的语义一一对应。Switch 在任务体系里属于操作符Operator类系统任务完整 JSON 配置说明可参考 switch-task.md。并行执行Fork / Join把互相独立的任务并行跑起来、等待全部完成后汇聚用ForkTaskJoinTaskfrom conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.workflow.task.fork_task import ForkTask from conductor.client.workflow.task.join_task import JoinTask workflow ConductorWorkflow(nameparallel_enrichment, version1, executorexecutor) # Define independent tasks credit_check check_credit(task_ref_namecredit, customer_idworkflow.input(customer_id)) fraud_check check_fraud(task_ref_namefraud, customer_idworkflow.input(customer_id)) kyc_check check_kyc(task_ref_namekyc, customer_idworkflow.input(customer_id)) # Fork runs all branches in parallel fork ForkTask( task_ref_nameparallel_checks, forked_tasks[ [credit_check], [fraud_check], [kyc_check], ], ) # Join waits for all branches join JoinTask(task_ref_namewait_all, join_on[credit, fraud, kyc]) # Merge results decide make_decision( task_ref_namedecide, credit_scorecredit_check.output(score), fraud_riskfraud_check.output(risk_level), kyc_statuskyc_check.output(status), ) workflow fork join decide workflow.output_parameters({decision: decide.output(result)}) workflow.register(overwriteTrue)要点拆解forked_tasks是分支列表的列表外层每个元素是一条并行分支内层列表是该分支要顺序执行的任务所以每个分支可以不止一个任务join_on声明 Join 等待哪些 taskReferenceName 完成。源码 JoinTaskMapper.java 显示Join 任务创建时把joinOn直接写入输入参数joinInput.put(joinOn, workflowTask.getJoinOn())服务端据此检查每个被等待分支是否达到终态Join 之后的任务示例中的decide可以同时引用多个分支的输出实现结果汇聚。若并行分支的数量在运行时才确定例如取决于 API 返回的条目数应使用 Dynamic ForkDYNAMIC_FORK_JOIN而非静态 Fork参考 dynamic-fork-task.md 与 dynamic-parallelism.md。循环Do/While当需要重复执行一组任务直到满足条件——例如轮询、重试、AI Agent 的思考-行动迭代——使用DoWhileTaskfrom conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.workflow.task.do_while_task import DoWhileTask workflow ConductorWorkflow(nameagent_loop, version1, executorexecutor) # The task(s) to repeat each iteration think call_llm( task_ref_namethink, promptworkflow.input(goal), ) act execute_tool( task_ref_nameact, toolthink.output(tool), argsthink.output(args), ) # Loop until the LLM says its done (max 10 iterations) 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}) summarize summarize_results(task_ref_namesummarize, resultsact.output(results)) workflow loop summarize workflow.register(overwriteTrue)要点拆解termination_condition是一段服务端求值的脚本表达式返回false表示继续循环返回true表示终止。示例中通过$.act[output][done]读取循环体内act任务的输出判断是否完成tasks是每次迭代要执行的任务列表max_iterations通过input_parameters.update(...)注入防止死循环。这也是循环控制的关键参数——Conductor 官方示例中有明确的 max 上限约定循环结束后act.output(results)仍可被循环后的summarize任务引用。服务端如何执行 DoWhile源码 DoWhileTaskMapper.java 展示了一个重要机制映射器先通过workflowModel.getTaskByRefName(...)查找同名循环任务如果该任务已处于终态terminal则直接返回空列表避免循环任务被重复调度。这意味着每次迭代都是服务端独立调度与持久化的任务执行天然具备失败重试与断点恢复能力——这正是动态循环适合承载 AI Agent 迭代的原因。完整配置参考 do-while-task.md。系统任务与自定义 Worker 混排HTTP、Wait、JSON JQ Transform 等属于系统任务system task由 Conductor 服务端 JVM 内建执行无需部署任何 Worker即可编排进工作流它们与worker_task自定义 Worker 可以在同一条链上自由组合from conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.workflow.task.http_task import HttpTask from conductor.client.workflow.task.json_jq_task import JsonJQTask from conductor.client.workflow.task.wait_task import WaitTask workflow ConductorWorkflow(namedata_pipeline, version1, executorexecutor) # HTTP task — fetch data from an external API (no worker needed) fetch HttpTask(task_ref_namefetch_data, http_input{ uri: https://api.example.com/records, method: GET, headers: {Authorization: [Bearer ${workflow.input.api_key}]}, }) # JQ Transform — reshape the response (no worker needed) transform JsonJQTask( task_ref_nametransform, script.body.records | map({id: .id, value: .metrics.total}), ) transform.input_parameters.update({ records: fetch.output(response.body), }) # Custom worker — run business logic enrich enrich_records( task_ref_nameenrich, recordstransform.output(result), ) # Wait — pause for 5 seconds before the next step cooldown WaitTask(task_ref_namecooldown, wait_for_seconds5) # Custom worker — store results store save_to_database(task_ref_namestore, recordsenrich.output(enriched)) workflow fetch transform enrich cooldown store workflow.output_parameters({stored: store.output(count)}) workflow.register(overwriteTrue)要点拆解HttpTask的http_input支持uri、method、headers、body等字段请求头里可以直接用${workflow.input.api_key}引用工作流入参注意${...}表达式在 Python 侧以字符串字面量传递由服务端解析JsonJQTask用script写 JQ 表达式重塑响应其输入records引用fetch.output(response.body)WaitTask(wait_for_seconds5)让工作流暂停 5 秒适用于冷却、节流、人工审批前的缓冲等场景系统任务在服务端 JVM 内执行、状态由服务端持久化因此整条链路依然是持久化执行durable execution——即使服务重启未完成的流程也能从最后一步恢复。Conductor 内置 20 系统任务与 10 种操作符Operators分类可参考 docs/devguide/concepts/tasks.md系统任务逐一配置说明见 systemtasks/index.md。子工作流把大流程拆成可复用组件当一个工作流体积过大时可拆成多个子工作流Sub Workflow——父工作流把子工作流当作一个普通任务调用。子工作流独立注册、独立版本化天然适合批处理每一条 父流程汇聚的模式from conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.workflow.task.sub_workflow_task import SubWorkflowTask # Child workflow (registered separately) child ConductorWorkflow(nameprocess_single_item, version1, executorexecutor) validate validate_item(task_ref_namevalidate, itemchild.input(item)) transform transform_item(task_ref_nametransform, itemvalidate.output(validated)) child validate transform child.output_parameters({result: transform.output(transformed)}) child.register(overwriteTrue) # Parent workflow invokes the child parent ConductorWorkflow(namebatch_processor, version1, executorexecutor) prepare prepare_batch(task_ref_nameprepare, batch_idparent.input(batch_id)) run_child SubWorkflowTask( task_ref_nameprocess_item, workflow_nameprocess_single_item, version1, ) run_child.input_parameters.update({item: prepare.output(first_item)}) aggregate aggregate_results( task_ref_nameaggregate, resultrun_child.output(result), ) parent prepare run_child aggregate parent.register(overwriteTrue)要点拆解子工作流child用child.input(item)声明自身入参与父工作流解耦父工作流用SubWorkflowTask(workflow_name..., version...)引用子工作流并通过input_parameters.update({item: ...})传入子工作流入参run_child.output(result)直接取子工作流输出——子工作流的output_parameters会透传到父级任务输出。从源码 SubWorkflowTaskMapper.java 可以看到两个值得注意的实现细节若没有指定子工作流版本映射器会通过MetadataDAO解析版本若传入的是内联定义inline workflow definition则直接使用内联定义而跳过版本解析映射器还支持priority、taskToDomain等附加参数通过subWorkflowTask.addInput(...)写入可用于子工作流的优先级与任务域隔离。完整配置参考 sub-workflow-task.md。运行时动态生成工作流免预注册即建即跑前面的模式都是先构建、再注册、后执行。Conductor 还支持第三种玩法——运行时现场构建工作流定义并直接启动执行无需预先注册。这对 AI Agent 特别有价值LLM 在运行时生成步骤清单代码把清单变成工作流定义Conductor 以持久化、可重试、可观测的方式执行它。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() # Build the workflow definition dynamically steps [validate, enrich, store] # determined at runtime 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}}, }, }) # Start with inline definition — no pre-registration needed 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})要点拆解steps列表就是运行时才知道的任务图——在上面的循环里每一步的输入都引用上一步的输出${validate_0.output.result}→${enrich_1.output.result}从而自动生成了链式依赖无需手写任何 JSON 字符串StartWorkflowRequest携带workflow_def内联定义与input工作流入参executor.start_workflow(request)立即返回workflow_id注意type: SIMPLE的语义这里指的是 Worker 任务类型需要对应名称的 Worker 在线轮询执行若希望完全不依赖 Worker可把类型换成HTTP、JSON_JQ_TRANSFORM等系统任务类型该模式与子工作流的内联定义能力一脉相承——SubWorkflowTaskMapper.java 中内联定义存在时跳过 MetadataDAO 版本解析的逻辑正是为了支持这类免注册的即建即跑场景。典型落地场景AI Agent 收到用户目标后LLM 生成一份执行计划步骤名、顺序、依赖关系你的代码将其翻译为tasks列表并即刻启动Conductor 负责后续的持久化执行、失败重试与全链路观测无需 Agent 自建状态机。同步执行并等待结果工作流也支持同步调用executor.execute(...)会阻塞到工作流完成直接返回运行结果适合 API 网关、交互式应用等请求-响应场景from conductor.client.configuration.configuration import Configuration from conductor.client.orkes_clients import OrkesClients config Configuration() clients OrkesClients(configurationconfig) executor clients.get_workflow_executor() # Execute synchronously — blocks until the workflow completes 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})execute与start_workflow的区别前者阻塞等待终态并返回结果对象后者立即返回workflow_id异步编排run.status、run.output直接可用config.ui_host拼接出的地址可直接在 Conductor UI 中打开该次执行的可视化详情页工作流执行会经历 RUNNING → COMPLETED / FAILED / TIMED_OUT / TERMINATED / PAUSED 等状态状态迁移图与timeoutPolicy、restartable等定义参数说明见 docs/devguide/concepts/workflows.md。环境准备与 SDK 安装所有示例都假定你已有一个WorkflowExecutor实例标准初始化如下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()pip install conductor-python export CONDUCTOR_SERVER_URLhttp://localhost:8080/api说明Configuration()默认从环境变量读取服务地址与认证信息CONDUCTOR_SERVER_URL、CONDUCTOR_AUTH_*本地单机部署时指向http://localhost:8080/api即可服务端部署方式见 docs/devguide/running/deploy.md让worker_task函数真正开始轮询任务需要启动TaskHandlerscan_for_annotated_workersTrue自动发现并为一个 Worker 函数启动一个轮询子进程完整可运行示例含 Worker 工作流 同步执行 UI 链接见 python-sdk.mdPython SDK 还支持异步 Workerasync def、长任务TaskInProgress、Prometheus 指标、工作流生命周期管理pause/resume/terminate/retry/restart/rerun/search等能力同样收录在该 SDK 文档中。小结编排能力Python SDK 构件服务端核心实现顺序执行操作符SIMPLEWorker 任务条件分支SwitchTask.switch_case / default_caseSwitchTaskMapper.javaEvaluator 求值 default 回退并行执行ForkTaskJoinTask(join_on...)JoinTaskMapper.java循环DoWhileTask(termination_condition, max_iterations)DoWhileTaskMapper.java终态去重系统任务HttpTask / JsonJQTask / WaitTask服务端 JVM 内建执行无需 Worker子工作流SubWorkflowTask(workflow_name, version)SubWorkflowTaskMapper.java支持内联定义运行时动态生成StartWorkflowRequest(workflow_def...)免预注册即建即跑代码优先的工作流让任务图在运行时才确定成为一等公民无论你的场景是订单履约、批处理管道还是由 LLM 驱动执行计划的 AI Agent都可以用同一套 Python 语法把它表达成持久化、可重试、可观测的 Conductor 工作流。文中所有模式对应的完整可运行示例可进一步参考仓库内 Python SDK 文档 python-sdk.md以及 first-workflow.md 中从零注册元数据到跑通首个工作流的实操路径。 /output文章 /output文章【免费下载链接】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),仅供参考