ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Conductor HTTP_POLL 任务实战:用单任务轮询长时间运行的外部作业

Conductor HTTP_POLL 任务实战:用单任务轮询长时间运行的外部作业 Conductor HTTP_POLL 任务实战用单任务轮询长时间运行的外部作业【免费下载链接】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导读当工作流需要向第三方 API 提交一个耗时几分钟甚至几小时的作业如批量导出、AI 推理、ETL 任务并等待其完成时传统做法是让工作流持有线程或维护一个工人进程成本高且脆弱。Conductor 内置的HTTP_POLL系统任务让你用一个任务完成提交 → 轮询状态 → 按结果分支的全流程它不占用线程、不需要自定义 worker也不会高频轰炸供应商接口。读完本文你将掌握terminationCondition、pollingInterval、pollingStrategy、maxPollCount四个轮询字段的完整用法能直接落地一份可运行的轮询工作流并理解其与DO_WHILE循环方案的优劣差异。场景一个慢作业一次等待你向某个第三方 API 提交了任务它返回一个jobId。这个作业要跑几分钟、有时几小时。你需要工作流等它跑完——但不能为此长期占用一个线程、不能部署一个专门轮询的 worker、也不能高频轰炸供应商的状态接口。HTTP_POLL正是为这个场景设计的一个系统任务你只需要给它状态 URL 和一个终止条件当这个条件为真就停剩下的轮询循环由 Conductor 服务器替你完成。整个模式的形状如下submit_job (HTTP) ── await_job (HTTP_POLL) ── SUCCEEDED ── record artifact │ polls the status URL FAILED ── TERMINATE │ until terminationCondition └─ sleeps between polls, holds nothing opensubmit_job用普通HTTP任务提交作业await_job用HTTP_POLL反复查询状态接口直到命中终止条件之后用SWITCH按最终状态分流成功则记录产物失败则用TERMINATE终止工作流。整个过程只有三个任务节点服务器在两次轮询之间睡眠不持有任何打开的连接或线程。为什么不用循环HTTP_POLL 与 DO_WHILE 对比用DO_WHILE包一个HTTP任务也能实现同样的效果你在旧示例里经常能看到这种写法。但它隐性的成本比表面上看起来高得多维度DO_WHILEHTTPHTTP_POLL执行中的任务数每次迭代两个无限增长一个轮询退避自己实现pollingStrategy轮询上限自己数迭代次数maxPollCount查看执行详情滚动翻过 40 次迭代一个任务 一个轮询计数更关键的是循环版本把最有趣的部分——终止条件——变成了埋在loopCondition里的表达式它只能基于循环状态求值而不是基于每次轮询的响应内容来求值。HTTP_POLL则把判断是否结束的职责交给terminationCondition该条件可以读取每次轮询返回的响应体、响应头和状态码语义直接、调试直观执行历史中只保留一个任务节点轮询次数作为任务本身的计数呈现。HTTP_POLL 任务定义在完整的工作流定义中HTTP_POLL任务的写法如下完整可运行版本见 docs/devguide/cookbook/assets/http-poll-external-job.json{ name: await_job, taskReferenceName: await_job, type: HTTP_POLL, inputParameters: { http_request: { uri: ${workflow.input.jobApiUrl}/jobs/${submit_job.output.response.body.jobId}, method: GET, terminationCondition: (function(){ var s $.output.response.body.state; return s SUCCEEDED || s FAILED; })();, pollingInterval: 60, pollingStrategy: FIXED, maxPollCount: 60 } } }HTTP_POLL接受与HTTP任务完全相同的http_request块——uri、method、headers、body、accept、contentType、connectionTimeOut、readTimeOut、acceptedStatusCodes、outputFilter——在此基础上再增加四个轮询专用字段字段默认值作用terminationCondition—每次轮询后求值的表达式返回真值即停止任务pollingInterval—两次轮询之间的间隔秒数pollingStrategy—FIXED、LINEAR_BACKOFF或EXPONENTIAL_BACKOFFmaxPollCount1000最多轮询多少次后放弃编写终止条件 terminationCondition终止条件表达式可以看到两个对象$.output—— 当前这次轮询的结果包括response.body、response.headers、response.statusCode$.input—— 任务自身的输入返回布尔值表示完成或继续。也可以返回数字实现三态控制1完成任务、0继续轮询、-1使任务失败。失败也要终止。一个只匹配SUCCEEDED的条件会在作业已经死掉的情况下继续空转轮询直到maxPollCount耗尽。应当匹配所有终态事后再分支处理结果(function(){ var s $.output.response.body.state; return s SUCCEEDED || s FAILED; })();这样一个条件同时覆盖成功与失败两个终态SUCCEEDED与FAILED之外的状态如QUEUED、RUNNING都会让任务继续轮询。轮询结束后下游用SWITCH按state分流详见下文可运行工作流定义一节失败的作业走TERMINATE分支终止工作流。轮询间隔存在服务器下限pollingInterval会被钳制到conductor.worker.http_poll.min_poll_interval该配置默认 60 秒。即使你请求pollingInterval: 5在运维没有调低下限的情况下实际生效的仍是 60 秒。因此maxPollCount要按生效间隔而不是你请求的间隔来估算60 次轮询 × 60 秒 一小时的上限若按 5 秒估算成 5 分钟上限实际任务会多跑近一小时才触发放弃逻辑。前置条件与桩服务要运行这个示例你需要一个正在运行的 Conductor 服务器以及一个可以被轮询的作业 API。仓库附带一个桩服务stub让你在没有供应商账号的情况下完整跑通整个流程。将下面的内容保存为job_stub_service.py并保持运行完整源码见 docs/devguide/cookbook/assets/job_stub_service.pyA stand-in for a slow third-party job API. Run it before starting the workflow: python3 job_stub_service.py # http://localhost:8089 Endpoints POST /jobs - 202, returns {jobId: ...} and starts a job GET /jobs/{id} - {jobId,state,progress,result} state goes QUEUED - RUNNING - SUCCEEDED POST /jobs/{id}/fail - force the job to FAILED on its next poll GET /polls - how many times each job has been polled The job advances one step per poll, so a workflow that polls it will see QUEUED, then RUNNING, then SUCCEEDED, without any wall-clock waiting. 启动它python3 job_stub_service.py # http://localhost:8089桩服务每次被轮询会让作业前进一个状态——QUEUED→RUNNING→RUNNING→SUCCEEDED——因此你无需等待真实的时间就能观察完整的生命周期。它提供的接口包括POST /jobs—— 返回202和{jobId: ..., state: QUEUED}创建一个新作业GET /jobs/{id}—— 返回{jobId, state, progress, result}每次轮询使状态前进一档POST /jobs/{id}/fail—— 强制作业在下一次轮询时变为FAILED用于演练失败分支GET /polls—— 返回每个作业被轮询的次数用于交叉验证轮询行为可运行的工作流定义将下面内容保存为http-poll-external-job.json并注册到 Conductor完整定义见 docs/devguide/cookbook/assets/http-poll-external-job.json{ name: http_poll_external_job, description: Submit a job to a slow third-party API, then let a single HTTP_POLL task poll it until it reports a terminal state. No loop task, no worker., version: 1, schemaVersion: 2, timeoutSeconds: 3600, timeoutPolicy: TIME_OUT_WF, inputParameters: [ jobApiUrl, dataset ], tasks: [ { name: submit_job, taskReferenceName: submit_job, type: HTTP, inputParameters: { http_request: { uri: ${workflow.input.jobApiUrl}/jobs, method: POST, body: { dataset: ${workflow.input.dataset} }, connectionTimeOut: 10000, readTimeOut: 30000 } } }, { name: await_job, taskReferenceName: await_job, type: HTTP_POLL, inputParameters: { http_request: { uri: ${workflow.input.jobApiUrl}/jobs/${submit_job.output.response.body.jobId}, method: GET, connectionTimeOut: 10000, readTimeOut: 30000, terminationCondition: (function(){ var s $.output.response.body.state; return s SUCCEEDED || s FAILED; })();, pollingInterval: 60, pollingStrategy: FIXED, maxPollCount: 60 } } }, { name: route_on_job_state, taskReferenceName: route_job, type: SWITCH, evaluatorType: value-param, expression: state, inputParameters: { state: ${await_job.output.response.body.state} }, decisionCases: { SUCCEEDED: [ { name: record_artifact, taskReferenceName: record_artifact, type: JSON_JQ_TRANSFORM, inputParameters: { body: ${await_job.output.response.body}, queryExpression: {jobId: .body.jobId, rows: (.body.result.rows // 0), artifact: (.body.result.artifact // \\)} } } ], FAILED: [ { name: terminate_job_failed, taskReferenceName: terminate_job_failed, type: TERMINATE, inputParameters: { terminationStatus: FAILED, workflowOutput: { error: remote_job_failed, jobId: ${await_job.output.response.body.jobId}, detail: ${await_job.output.response.body.error} } } } ] }, defaultCase: [] } ], outputParameters: { jobId: ${submit_job.output.response.body.jobId}, finalState: ${await_job.output.response.body.state}, rows: ${record_artifact.output.result.rows}, artifact: ${record_artifact.output.result.artifact} } }这个定义里有几个值得注意的工程细节submit_job的响应被HTTP_POLL复用await_job的状态 URL 通过${submit_job.output.response.body.jobId}引用提交任务的返回体任务间通过输出引用传递数据无需额外状态存储。SWITCH用value-param求值器直接按await_job响应体中的state字段分流。SUCCEEDED分支用JSON_JQ_TRANSFORM从响应中提取jobId、rows、artifact并结构化记录FAILED分支用TERMINATE以FAILED状态结束工作流并把remote_job_failed与作业详情写进工作流输出。工作流级兜底timeoutSeconds: 3600与timeoutPolicy: TIME_OUT_WF为整个执行设定一小时上限防止轮询配置失误时工作流无限悬挂。注册与运行使用 Conductor CLI 注册工作流定义并启动一次执行conductor workflow create http-poll-external-job.json conductor workflow start -w http_poll_external_job \ -i {jobApiUrl:http://localhost:8089,dataset:orders_2026_q2}然后打开 Conductor UI 的Executions页面选中新生成的执行查看任务图以及每个任务的输入与输出。运行过程中你可以观察到await_job始终是一个任务它的轮询计数不断攀升当桩服务返回SUCCEEDED时SWITCH走成功分支记录产物。若要演练失败分支先启动工作流再用桩服务强制失败curl -X POST http://localhost:8089/jobs/{jobId}/fail之后的工作流会以remote_job_failed终止而不是继续轮询一个已死的作业。交叉验证供应商实际看到的轮询行为curl -s http://localhost:8089/polls返回的polls映射能直接对照await_job的轮询计数确认两次轮询之间确实隔着pollingInterval秒。生产环境注意事项把上面的示例搬到生产环境时以下几条原则直接决定可靠性终止条件必须覆盖所有终态而不只是成功。否则一个已经死掉的作业会让任务一直轮询到maxPollCount耗尽白白消耗配额与时间。pollingInterval有服务器端下限min_poll_interval默认 60 秒。你配置的值只是一个请求不是保证最终生效值以服务器钳制后的结果为准。用墙钟时间预算来设置maxPollCount。间隔 × 次数才是真实上限同时给工作流设置一个略高于该上限的timeoutSeconds作为兜底避免轮询配置错误导致执行无限悬挂。作业时长未知时优先用EXPONENTIAL_BACKOFF。一个跑 5 小时的作业若用固定间隔会产生 300 次相同的请求指数退避让前段轮询密集、后段稀疏显著降低请求量。轮询要打便宜接口。如果供应商的状态接口被限流或返回完整负载尽量申请一个轻量级的状态 URL或用outputFilter过滤响应避免把大响应体写进工作流状态。提交步骤需要幂等键。被重试的提交可能创建出第二个作业导致你轮询错了对象为提交请求设计幂等机制如携带唯一requestId以规避。不要用于亚秒级工作。在轮询下限之下同步HTTP任务才是正确工具——HTTP_POLL的粒度由服务器下限决定。关联阅读等待与定时器模式 —— 相比轮询状态 URL如何在信号或时钟上等待任务超时与重试 —— 为提交调用设置边界Saga补偿部分失败 —— 当后续步骤失败时如何撤销已提交的作业【免费下载链接】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),仅供参考
返回列表