Conductor 如何用 Python SDK 以代码方式(workflow as code)动态构建工作流
2026/9/10 22:11:24 网站建设 项目流程

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/conductor

Conductor 支持 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_URL=http://localhost:8080/api

Python 中获取执行器的标准连接代码:

from conductor.client.configuration.configuration import Configuration from conductor.client.orkes_clients import OrkesClients config = Configuration() # reads CONDUCTOR_SERVER_URL from env clients = OrkesClients(configuration=config) executor = clients.get_workflow_executor()

一个容易遗漏的前提:工作流中的自定义任务(SIMPLE 任务)需要 worker 轮询执行。SDK 提供TaskHandler,它会自动发现所有@worker_task装饰的函数并为每个 worker 启动一个子进程:

from conductor.client.automator.task_handler import TaskHandler with TaskHandler(configuration=config, scan_for_annotated_workers=True) as task_handler: task_handler.start_processes()

如果 worker 没有在跑,工作流会启动但任务一直处于等待状态,同步执行会一直阻塞。

用 >> 运算符构建顺序工作流

@worker_task装饰的普通 Python 函数就是可复用的任务积木,task_definition_name是任务的注册名。文档中的示例(订单履约流程)如下,其中的返回值(如99.99txn_abc123)是文档示例值,实际业务请替换函数体:

from conductor.client.workflow.conductor_workflow import ConductorWorkflow from conductor.client.worker.worker_task import worker_task @worker_task(task_definition_name='fetch_order') def fetch_order(order_id: str) -> dict: return {'order_id': order_id, 'amount': 99.99, 'item': 'Widget'} @worker_task(task_definition_name='process_payment') def process_payment(order_id: str, amount: float) -> dict: return {'transaction_id': 'txn_abc123', 'status': 'charged'} @worker_task(task_definition_name='ship_order') def ship_order(order_id: str, transaction_id: str) -> dict: return {'tracking': 'TRACK-456', 'carrier': 'FedEx'} workflow = ConductorWorkflow(name='order_fulfillment', version=1, executor=executor) fetch = fetch_order(task_ref_name='fetch', order_id=workflow.input('order_id')) pay = process_payment( task_ref_name='pay', order_id=workflow.input('order_id'), amount=fetch.output('amount'), ) ship = ship_order( task_ref_name='ship', order_id=workflow.input('order_id'), transaction_id=pay.output('transaction_id'), ) workflow >> fetch >> pay >> ship workflow.output_parameters({ 'tracking': ship.output('tracking'), 'transaction_id': pay.output('transaction_id'), }) workflow.register(overwrite=True)

几个关键用法:

  • workflow.input('order_id'):引用工作流启动时的输入字段,作为任务输入参数;
  • fetch.output('amount'):引用上游任务输出中的某个字段,作为下游任务的输入,任务间的数据传递就这样串起来;
  • task_ref_name是任务在流程内的引用名;
  • register(overwrite=True)把工作流定义注册到服务器,overwrite=True表示已存在同名版本时覆盖。

把 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_name='greet', register_task_def=True) def greet(name: str) -> str: return f'Hello {name}' def main(): config = Configuration() clients = OrkesClients(configuration=config) executor = clients.get_workflow_executor() workflow = ConductorWorkflow(name='greetings', version=1, executor=executor) greet_task = greet(task_ref_name='greet_ref', name=workflow.input('name')) workflow >> greet_task workflow.output_parameters({'result': greet_task.output('result')}) workflow.register(overwrite=True) # 启动 worker 子进程开始轮询 with TaskHandler(configuration=config, scan_for_annotated_workers=True) as task_handler: task_handler.start_processes() run = executor.execute(name='greetings', version=1, workflow_input={'name': 'Conductor'}) print(f'result: {run.output["result"]}') print(f'execution: {config.ui_host}/execution/{run.workflow_id}') if __name__ == '__main__': main()

注意register_task_def=True的用途:它让 SDK 在本地开发时顺带注册任务定义,文档同时提示生产环境中应单独管理任务定义,不依赖这个开关。

同步执行并验证结果

executor.execute会阻塞直到工作流完成,返回对象的statusoutput就是文档给出的验证方式:

run = executor.execute( name='order_fulfillment', version=1, workflow_input={'order_id': 'ORD-789'}, ) print(f'Status: {run.status}') print(f'Output: {run.output}') print(f'View: {config.ui_host}/execution/{run.workflow_id}')

判断是否成功的两条路径:

  1. 打印的run.output中应包含output_parameters声明的键(上面示例即trackingtransaction_id),值对应各 worker 的返回值;
  2. 打开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(configuration=config) 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( name='dynamic_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(f'Started dynamic workflow: {workflow_id}')

两点说明:

  • 代码里的${workflow.input.data}${validate_0.output.result}等是 Conductor 的输入表达式语法(写在任务inputParameters里,由服务器在执行时解析),不是 Python 变量或需要人工替换的模板;Python f-string 负责在运行时把它们拼成正确的字符串;
  • 这条路径走start_workflow异步启动,立即返回workflow_id。验证方式是拿这个 ID 去 UI 检查执行状态;validateenrichstore对应的 worker 必须已在运行,任务才会被消费。

条件分支、并行与循环

以下三种结构用于替代手工 JSON 时的等价能力,示例中的 worker 函数(classify_ticketpage_oncallcheck_credit等)文档未给出函数体,需要你用@worker_task自行定义后再套用。

Switch 条件分支——按任务输出路由,每个 case 是一条独立的任务链:

from conductor.client.workflow.task.switch_task import SwitchTask switch = SwitchTask(task_ref_name='priority_router', case_expression=classify.output('priority')) switch.switch_case('critical', [ page_oncall(task_ref_name='page', ticket_id=workflow.input('ticket_id')), escalate(task_ref_name='escalate', ticket_id=workflow.input('ticket_id')), ]) switch.switch_case('high', [ assign_senior(task_ref_name='assign', ticket_id=workflow.input('ticket_id')), ]) switch.default_case([ add_to_backlog(task_ref_name='backlog', ticket_id=workflow.input('ticket_id')), ]) workflow >> classify >> switch

Fork/Join 并行——ForkTaskforked_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_name='parallel_checks', forked_tasks=[[credit_check], [fraud_check], [kyc_check]], ) join = JoinTask(task_ref_name='wait_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_name='agent_loop', termination_condition='if ($.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(overwrite=True)会覆盖同名版本的既有定义,注意不要用它覆盖线上正在使用的版本;
  • 完整生命周期操作(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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询