完整指南:从 @workflow 到人工审批)
Windmill Python Workflow-as-Code APIwmill完整指南从 workflow 到人工审批【免费下载链接】windmillOpen-source developer platform to power your entire infra and turn scripts into webhooks, workflows and UIs. Fastest workflow engine (13x vs Airflow). Open-source alternative to Retool and Temporal.项目地址: https://gitcode.com/GitHub_Trending/wi/windmillWindmill 的 Python Workflow-as-CodeWACAPI 让你用纯 Python 代码声明工作流workflow装饰异步入口函数task把函数拆成独立子任务step/sleep/wait_for_approval负责轻量内联计算、服务端睡眠与人工审批parallel做并发控制。本文以仓库中 system_prompts/auto-generated/sdks/wac-python.md 为骨架结合 python-client/wmill/wmill/client.py 与 python-client/wmill/tests/test_workflow.py 的源码实现完整讲解每个 API 的签名、参数、行为与底层原理读完即可直接用 wmill 编写可落地的 WAC 工作流。一、先看全貌WAC 的两种执行模式WACWorkflow-as-Code是 Windmill 把代码即工作流落到 Python 的方式。task装饰的函数在三种上下文中表现不同源码见 client.pyWAC v2在workflow内部任务被调度为检查点步骤checkpoint step父作业挂起子任务作为独立作业运行结果写回检查点后父作业重放replay继续WAC v1设置了WM_JOB_ID但不在workflow中任务通过 HTTP API 同步派发对应源码中的/w/{workspace}/jobs/run/workflow_as_code/{job_id}/{func_name}端点client.py独立运行本地/脚本直接执行函数体但结果仍会做一次 JSON 往返round-trip保证本地跑出的类型和部署后一致。workflow装饰的函数必须确定性deterministic同样的输入每次重放必须以相同顺序调用任务。基于任务结果分支是允许的结果从检查点重放但基于外部状态当前时间、随机值、外部 API 调用分支则必须用step()把值先写入检查点让重放看到相同结果client.py。顶层导入语句与文档一致from wmill import workflow, task, task_script, task_flow, step, sleep, wait_for_approval, get_approval_urls, get_resume_urls, parallel, TaskError二、异常处理TaskError当 WAC 的task或step失败时抛出TaskError继承自Exception构造签名与属性见 client.pyclass TaskError(Exception): def __init__(self, message: str, *, step_key: str , child_job_id: Optional[str] None, result None)三个关键属性step_key失败步骤的检查点键checkpoint keychild_job_id失败子作业的 UUID对step()为None因为step()运行在工作流作业内部没有子作业result统一形状{error: {name, message, stack?, extra?}}——无论 task 还是 step 失败形状一致。name和message恒有stack仅在失败带 traceback 时出现extra仅在异常携带自定义字段时出现若字段太大无法写入检查点会被丢弃并附带extra_omitted: True。源码中失败的step()体通过_step_error_marker序列化成__wmill_error标记写入completed_steps重放时由_task_error_from_marker重建为TaskError。值得注意的设计产生失败的那一轮和后续每一轮重放都走同一个重建函数client.py——因为workflow每轮都从函数顶部重新执行except分支是控制流必须保证每轮抛出完全相同的异常对象否则根据失败分支的重放会派发不同的任务。三、核心任务原语task 与任务选项task装饰器完整签名def task(_func None, *, path: Optional[str] None, tag: Optional[str] None, timeout: Optional[int] None, cache_ttl: Optional[int] None, priority: Optional[int] None, concurrency_limit: Optional[int] None, concurrency_key: Optional[str] None, concurrency_time_window_s: Optional[int] None, retry: Optional[dict] None)典型用法文档示例均可用task async def extract_data(url: str): ... task(pathf/external_script, timeout600, taggpu) async def run_external(x: int): ... task(retry{attempts: 3, delay: 30, multiplier: 2}) async def call_api(payload: dict): ...要点任务作为独立作业运行结果总是先 JSON 编码再解码回来——datetime会变成字符串、tuple会变成list源码中_json_round_trip就是模拟这一行为见 client.py支持task无括号与task(...)带参数两种写法client.pypath指定后任务派发到另一个 Windmill 脚本而非函数本身此时task_name取函数名、script取pathclient.py。retry 重试策略详解retry仅在workflow内部生效接受一个字典合法键固定为(attempts, delay, multiplier, max_delay)源码常量_RETRY_KEYSclient.py键含义约束attempts首次失败后的重试次数必填0–100 的整数delay首次重试前的等待秒数亚秒级延迟会被丢弃小于 1 秒视为 0multiplier每次尝试后对延迟的乘数1 表示恒定缺省为 1max_delay延迟上限秒可选策略在写入处即被校验_checked_retryclient.py未知键立即抛ValueErrorattempts非 0–100 的整数含布尔值也会被拒绝——因为策略是普通 dict拼错的键若静默丢弃任务就会按没人写过的策略重试所以必须在装饰时拒之门外。实现层面的两个关键设计每次尝试都是独立步骤第 1 次尝试的检查点键是call_api重试依次为call_api#2、call_api#3……两次尝试之间的等待是持久化睡眠durable sleep所以正在退避backoff的重试任务不占用任何 workerclient.py。_retry_delay_seconds计算每次退避时长溢出时封顶到_MAX_SLEEP_SECONDS 2**32 - 1worker 端把睡眠反序列化为 u32 秒数超宽会整个作业失败见 client.py扇出fan-out中的退避是串行求和而非并行取最长工作流每轮只睡一次同一扇出中正在退避的任务会一个接一个地等扇出重试前的总延迟等于所有待退避延迟之和随扇出宽度与attempts增长不带delay的重试全部在同一轮发出。重试键在第一次尝试派发前就被一次性认领完因此无界的attempts不会导致工作流挂死分配而是在写策略处就被上限 100 拦下。四、复用已有资产task_script 与 task_flow如果你的任务逻辑已经以 Windmill 脚本或 Flow 的形式存在不必重写# 派发到独立脚本 extract task_script(f/data/extract, timeout600) # 派发到独立 Flow pipeline task_flow(f/etl/pipeline, priority10) workflow async def main(): data await extract(urlhttps://...) result await pipeline(inputdata)二者签名一致path必填其余timeout、tag、cache_ttl、priority、concurrency_limit、concurrency_key、concurrency_time_window_s、retry均可选retry策略与task完全相同task_script内部以dispatch_typescript、task_flow以dispatch_typeflow调用ctx._next_stepclient.py、client.py二者只能在工作流内部调用在工作流外调用会抛RuntimeError。五、轻量内联计算stepstep在父作业内内联执行函数并把结果写入检查点async def step(name: str, fn)重放时直接返回缓存值不再执行fn适用于轻量、确定性的操作——时间戳、随机 ID、配置读取等不值得开一个子作业的场景fn的结果同样经 JSON 编码/解码因此跑函数体的那一轮看到的类型和每一轮重放一致datetime变字符串、tuple变list。源码中step走WorkflowCtx._run_inline_stepclient.py默认开启内联快速路径由环境变量WM_WAC_INLINE_FAST_PATH控制默认1设为0/false/off/no可回退到旧的挂起-重放路径直接把结果 POST 到/api/w/{workspace}/jobs/wac/inline_checkpoint/{job_id}端点而不必展开协程栈失败时退化为传统的 suspend-and-replay 路径。并发step()之间只串行化 HTTP 请求通过_asyncio.Lockfn()本身仍可并行执行——asyncio.gather(step(a, fn_a), step(b, fn_b))中两个函数体并行只有 API 写入按序client.py。六、服务端睡眠sleepasync def sleep(seconds: int)在workflow内部父作业挂起seconds秒后自动恢复全程不占用 worker源码中转为mode: sleep的_StepSuspend见 client.py在工作流之外退化为asyncio.sleep。检查点键自动分配sleep、sleep_2……若键已在检查点中则直接跳过等待。七、并行与并发控制parallelasync def parallel(items, fn, *, concurrency: Optional[int] None)fn应当是一个task每个 item 通过fn(item)处理默认一次性全部派发concurrencyNone时批量大小等于len(items)指定concurrency时按批次派发例如concurrency5每批 5 个空列表直接返回[]。源码实现client.py按批使用asyncio.gather并发执行任务在 WAC v2 下同一轮收集到的多个待派发任务会被打包成mode: parallel的一次挂起_suspend中len(steps) 1判定为 parallel见 client.py。八、人工审批wait_for_approval 与审批 URLwait_for_approvalasync def wait_for_approval(timeout: int 1800, form: dict | None None, self_approval: bool True, key: str | None None, skin: Literal[detailed, minimal] | None None, description: str | dict | None None) - dict挂起工作流等待外部审批返回包含value表单数据、approver、approved的字典。参数说明参数说明timeout审批超时秒数默认 180030 分钟form审批页的可选表单 schemaself_approval触发该流程的用户能否自行审批默认Truekey可选检查点键为本次审批步骤命名skinminimal时审批者只看到请求表单与批准/拒绝而非带工作流详情的大页面description显示在表单上方给审批者的说明字符串或富值如{markdown: ...}不传key时步骤自动命名为approval、approval_2……源码_alloc_key的按名分配机制client.py。显式传入key时键必须在工作流内唯一——重复使用会抛RuntimeError而不是静默改名因为 URL 是按键铸造的静默改名会让调用者拿到指向第一个步骤的 URL最终报 resume request already sent 并把工作流挂到超时client.py。文档给出的完整示例urls await step(urls, lambda: get_approval_urls(manager)) await step(notify, lambda: send_email(urls[resume], urls[cancel])) result await wait_for_approval(keymanager, timeout3600)get_approval_urls 与 get_resume_urls两个函数都返回含approvalPage、resume、cancel三个 URL 的字典但语义不同def get_approval_urls(step_key: str approval, approver: str None) - dict def get_resume_urls(approver: str None, flow_level: bool None) - dictget_approval_urls(step_key, ...)绑定到某一个wait_for_approval步骤与get_resume_urls不同它不是用随机 nonce 签名而是直接定位该步骤内置审批按钮所用的那条resume_job记录因此跨重放稳定可安全嵌入自定义通知client.pyresume与cancel是步骤绑定的只在步骤等待审批期间可用其他时刻使用会被拒绝而不是预存一行给别的审批消费可以提前发送但审批者在工作流到达该步骤前无法操作approvalPage不是步骤绑定的它打开作业的审批页作用于使用时正在等待中的那个审批源码中请求GET /w/{workspace}/jobs/wac_approval_urls/{job_id}/{quote(step_key)}client.pystep_key经 URL 编码。get_resume_urls(approverNone, flow_levelNone)面向传统挂起suspend步骤approver可选审批者名称flow_levelTrue时为父流程而非具体步骤生成恢复 URL从而支持预审批——同一流程中任何后续挂起步骤都可消费该审批源码中通过随机 nonce 调用GET /w/{workspace}/jobs/resume_urls/{job_id}/{nonce}client.py。九、底层原理检查点键分配与工作流重放理解 WAC v2 的关键是WorkflowCtxclient.py的检查点模型按名分配键同一名称首次调用得double后续得double_2、double_3……且会跳过已被占用的键避免step(x)第二次调用与第一次step(x_2)冲突。分配顺序由工作流函数体决定因此重放必然复现相同键序列每轮只派发一次工作流体执行到某个任务/步骤时把派发信息压入_pending并抛_StepSuspend_run_workflow_asyncclient.py捕获后按mode返回dispatch/inline_checkpoint/approval/sleep/complete等指令给 workerworker 执行完成后带着新检查点重新启动工作流体重放时命中completed_steps的键直接返回缓存值重试键提前认领带retry的任务在首次派发前就把base_key、base_key#retry2、base_key#2等全部键分配好避免稍后分配导致旁边步骤的键偏移client.py。测试侧test_workflow.py 覆盖了WorkflowCtx、_StepSuspend、TaskError、workflow、task、step、sleep、parallel、wait_for_approval等全部核心原语并包含_StubInlineClient验证内联快速路径 POST 到/jobs/wac/inline_checkpoint的请求体可作为理解各 API 行为边界的参考。十、快速上手指南在 Windmill 中启用 Python WAC 的典型步骤安装客户端项目依赖声明于 python-client/wmill/pyproject.toml运行环境由 Windmill 的 Python worker 自动注入编写工作流顶层导入from wmill import workflow, task, step, sleep, parallel, wait_for_approval, ...用workflow装饰async def main()内部用task函数与step/sleep/parallel/wait_for_approval编排运行环境变量客户端通过WM_TOKEN、WM_WORKSPACE、WM_BASE_URL或BASE_INTERNAL_URL、WM_JOB_ID等环境变量自动感知部署上下文client.py独立运行时可用WM_MOCKED_API_FILE指向 mock API 文件做离线联调本地验证在工作流外调用task/step会直接执行函数体并做 JSON 往返保证本地结果与部署一致。小结Windmill 的 Python WAC API 把代码即工作流落到一个十来个函数的 Python 接口上workflowtask构成骨架step/sleep/parallel提供内联计算、服务端睡眠与并发控制wait_for_approvalget_approval_urls/get_resume_urls打通人工审批TaskError统一失败语义。底层是确定性的检查点/重放模型——每轮执行到挂起点即返回worker 执行子任务后携带新检查点重放函数体由此获得无需占住 worker 的持久睡眠、跨重放稳定的重试退避与审批 URL。理解这一模型就能写出既能在本地快速调试、又能在 Windmill 集群上可靠长跑的 WAC 工作流。【免费下载链接】windmillOpen-source developer platform to power your entire infra and turn scripts into webhooks, workflows and UIs. Fastest workflow engine (13x vs Airflow). Open-source alternative to Retool and Temporal.项目地址: https://gitcode.com/GitHub_Trending/wi/windmill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考