ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

零成本搭建AI工作流:编排与执行分离的工程实践

零成本搭建AI工作流:编排与执行分离的工程实践 1. 六毛钱引发的思考为什么我要自己搭一套AI工作流事情的起因特别简单。上个月我在整理一批资料需要把几十份文档做摘要、分类、提取关键信息再统一转成固定格式输出。我一开始用的是某个在线AI工具单次处理收费六毛钱。听起来不多对吧但我那天要跑两百多次算下来一百多块就没了。更让我难受的是这个工具的处理逻辑是固定的我没法调整中间步骤也没法把结果直接接到我自己的后续流程里。我当时就想这东西的本质不就是“把文本喂给模型拿回结果再做后处理”吗这套逻辑我自己也能搭而且搭完之后是零边际成本的——跑一次和跑一万次除了电费和模型调用费不再有额外的工具订阅开销。于是就有了这套零成本AI工作流的设计。先把话说在前面这里的“零成本”不是指完全不花钱。模型API调用本身是有费用的如果你用的是本地模型那电费和硬件折旧也是成本。我说的零成本是指不再为工作流工具本身付费不再被某个平台的订阅制绑架整个流程的编排、调度、存储、后处理全部用免费或开源方案搞定。对于每天要处理几十上百条任务的个人用户和小团队来说这个差别非常大。这套工作流适合什么人我总结了三类第一类是像我这样有批量文本处理需求的个人比如做内容整理、资料归档、数据清洗第二类是想把AI能力接进自己现有系统但不想被平台锁死的开发者第三类是单纯想搞明白“AI工作流到底是怎么回事”的学习者。如果你属于这三类中的任何一类接下来的内容应该对你有用。在正式开始拆解之前我需要先明确一个核心设计原则这个原则贯穿整套工作流的始终把“编排”和“执行”彻底分开。编排层负责决定“先做什么、再做什么、什么条件下走哪个分支”执行层负责“真正去调用模型、读写文件、发请求”。这样设计的好处是编排层可以随时改执行层可以随时换两边互不影响。很多平台级工具之所以让人用着难受就是因为它们把这两层揉在一起了你想改一个判断条件结果发现得把整个流程重新配一遍。2. 拆解工作流的核心骨架编排层与执行层到底怎么分2.1 编排层用配置文件描述“做什么”编排层我选择用YAML配置文件来描述整个流程。为什么不用可视化拖拽因为可视化拖拽在流程简单的时候很直观但一旦分支变多、条件变复杂那堆连线和方框反而比代码更难维护。YAML的好处是纯文本、可版本管理、可复用、可参数化而且改起来就是改几行字的事。一个典型的编排配置长这样workflow: name: document_processor steps: - id: extract_text type: file_read params: path: {{input_path}} encoding: utf-8 - id: summarize type: llm_call params: model: local_model prompt: 请对以下内容做200字以内的摘要\n{{extract_text.output}} max_tokens: 500 - id: classify type: llm_call params: model: local_model prompt: 请判断以下摘要属于哪个类别技术/生活/商业/其他\n{{summarize.output}} - id: save_result type: file_write params: path: {{output_dir}}/{{extract_text.filename}}_result.md content: ## 摘要\n{{summarize.output}}\n\n## 分类\n{{classify.output}}这个配置里每个步骤都有id、type和params。type决定了这一步由哪个执行器来处理params里的{{...}}是变量引用指向前面步骤的输出。整个流程的依赖关系通过变量引用来表达不需要额外画箭头。提示变量引用的设计很关键。我一开始用的是“步骤A的输出传给步骤B”这种显式连线方式后来发现一旦步骤多了改一个中间步骤就得把所有下游引用全改一遍。改成{{step_id.output}}这种命名引用之后增删步骤只需要改引用它的地方维护成本低了很多。2.2 执行层每个type对应一个执行器执行层是一组独立的执行器函数每个执行器只负责一件事。比如file_read执行器只负责读文件llm_call执行器只负责调模型file_write执行器只负责写文件。执行器之间不直接通信所有数据传递都通过编排层统一管理。这种设计有个很实际的好处你可以单独测试任何一个执行器。比如你想验证llm_call执行器在某个模型上的表现不需要跑整个流程直接给它喂一段文本就行。我在调试阶段就是这么干的先把每个执行器单独跑通再串起来跑完整流程排查问题的效率高很多。执行器的接口我统一成了这样def execute(params: dict, context: dict) - dict: params: 编排层传入的参数已经解析过变量引用 context: 全局上下文包含输入路径、输出目录、配置等 返回: {output: ..., filename: ..., status: success} 所有执行器都遵循这个签名新增一个执行器只需要写一个函数然后注册到类型映射表里。我目前注册了六种类型file_read、file_write、llm_call、http_request、text_transform、condition。对于大多数文本处理场景这六种已经够用了。2.3 变量解析整个工作流的“神经系统”变量解析是编排层最核心的机制。当执行器拿到params时里面的{{...}}已经被替换成了实际值。解析逻辑分三步先扫描所有{{...}}表达式再根据表达式去上下文里取值最后把值替换回原字符串。取值的时候支持两种路径一种是{{step_id.output}}直接取某个步骤的输出另一种是{{step_id.output.field}}取输出字典里的某个字段。如果路径不存在解析器会抛出一个明确的错误告诉你哪个步骤的哪个字段找不到而不是静默返回空值。这一点很重要静默失败是工作流调试中最让人头疼的问题。注意变量解析的顺序必须严格按照步骤的依赖关系来。如果步骤B引用了步骤A的输出那A必须先执行完。我在实现的时候用了一个简单的拓扑排序根据变量引用关系确定执行顺序。如果你的流程里有循环引用排序会直接报错不会陷入死循环。3. 模型调用这一环本地模型和API怎么选、怎么接3.1 本地模型零边际成本的核心既然目标是零成本本地模型是绕不开的一环。我目前在用的是一台带独立显卡的台式机跑一个7B参数量的量化模型处理摘要和分类这类任务完全够用。本地模型最大的好处是没有调用次数限制你想跑多少遍就跑多少遍不用担心账单。但本地模型也有明显的短板第一首次加载慢模型越大越明显第二长文本处理能力有限上下文窗口通常比云端模型小第三复杂推理任务的表现不如大参数模型。我的应对策略是分层处理简单任务摘要、分类、格式转换走本地模型复杂任务长文分析、多步推理走云端API。这样既控制了成本又保证了效果。本地模型的接入方式我用的是兼容OpenAI接口的本地服务。这样做的好处是执行器里的llm_call不需要区分“本地”还是“云端”只需要改base_url和model两个参数就行。切换模型的时候编排配置里改一行执行器代码完全不用动。# llm_call 执行器的核心逻辑 import requests def execute(params, context): base_url params.get(base_url, context[default_base_url]) model params.get(model, context[default_model]) prompt params[prompt] response requests.post( f{base_url}/v1/chat/completions, json{ model: model, messages: [{role: user, content: prompt}], max_tokens: params.get(max_tokens, 1000), temperature: params.get(temperature, 0.3) }, timeout120 ) result response.json() return {output: result[choices][0][message][content]}3.2 云端API什么时候值得花这个钱云端API我只在两种情况下用一是本地模型确实搞不定的复杂任务二是对延迟敏感、需要快速返回的场景。选API的时候我主要看三个指标每百万token的价格、上下文窗口大小、以及是否支持流式输出。这里有个经验不要只看单价要看“有效单价”。有些模型单价便宜但输出质量差你需要反复重试或者人工修正实际成本反而更高。我实测下来对于摘要和分类任务一个中等价位的模型加上好的提示词比一个便宜模型加上反复重试要划算得多。还有一个容易被忽略的点API调用的超时和重试策略。我一开始没做重试结果网络抖动一下整个流程就断了。后来加了指数退避重试最多重试三次每次间隔翻倍。这个改动让流程的稳定性提升了一个档次。import time def call_with_retry(url, payload, max_retries3): for attempt in range(max_retries): try: resp requests.post(url, jsonpayload, timeout120) if resp.status_code 200: return resp.json() except requests.exceptions.RequestException: pass if attempt max_retries - 1: time.sleep(2 ** attempt) raise RuntimeError(fAPI调用失败已重试{max_retries}次)3.3 提示词模板让模型输出稳定的关键工作流里的模型调用和平时聊天不一样你需要的是稳定、可预测的输出。我踩过的最大坑就是提示词写得太随意导致同一个输入跑两次得到两种格式的输出后处理直接崩掉。我的解决方案是结构化提示词模板。每个llm_call步骤的提示词都遵循“角色任务格式要求示例”的四段式结构。比如摘要任务的提示词你是一个文档摘要助手。请对以下内容做摘要要求 1. 摘要长度在150-250字之间 2. 保留关键数据和结论 3. 不要添加原文没有的信息 4. 直接输出摘要正文不要加“摘要”前缀 内容 {{extract_text.output}}格式要求越具体输出越稳定。如果任务对格式要求特别严格我还会在提示词里加一个输出示例让模型照着抄格式。实测下来加了示例之后格式错误率从大概两成降到了几乎为零。4. 从输入到输出一条完整链路的实操拆解4.1 输入处理文件读取与预处理输入环节看起来简单其实有不少细节。我处理的文件格式主要是Markdown、纯文本和PDF转出来的文本。Markdown和纯文本直接读就行PDF转文本需要额外处理因为转换工具经常会带出多余的换行和页眉页脚。我的做法是在file_read执行器里加一个可选的clean参数。如果clean: true读取之后会自动做几件事合并连续空行、去掉每行首尾空白、去掉常见的页眉页脚模式比如“第X页”。这些清洗规则不复杂但能显著提升后续模型处理的质量。import re def clean_text(text): # 合并连续空行 text re.sub(r\n{3,}, \n\n, text) # 去掉行首行尾空白 lines [line.strip() for line in text.split(\n)] # 去掉页码行 lines [line for line in lines if not re.match(r^第?\s*\d\s*页?$, line)] return \n.join(lines)还有一个实际问题是文件编码。我遇到过好几次UTF-8和GBK混用的情况读出来全是乱码。后来在file_read里加了编码探测先试UTF-8失败再试GBK再失败就用errorsreplace兜底。虽然不能保证百分百正确但至少不会让整个流程崩掉。4.2 中间处理多步调用的串联与数据传递中间处理是工作流的核心价值所在。单个模型调用谁都会但把多个调用串起来、让上一步的输出成为下一步的输入、并且保证数据格式一致这才是工作流真正解决的问题。我举一个实际的例子。我有一批技术文章需要处理流程是这样的先提取正文然后生成摘要再根据摘要判断文章难度等级最后按难度等级输出到不同目录。四个步骤三个模型调用一个文件写入。这里的关键是步骤之间的数据契约。summarize步骤的输出必须是纯文本摘要不能带任何额外格式classify步骤的输出必须是“初级/中级/高级”三个词之一不能是“这篇文章属于中级难度”这种句子。为了保证这一点我在提示词里明确要求“只输出一个词”并且在执行器里加了一个输出校验如果输出不在预期范围内就重试一次重试还不行就标记为“待人工处理”。VALID_LEVELS {初级, 中级, 高级} def validate_classify(output): cleaned output.strip() if cleaned in VALID_LEVELS: return cleaned # 尝试从输出中提取 for level in VALID_LEVELS: if level in cleaned: return level return None # 标记为待人工处理这个校验机制看起来不起眼但它把“模型偶尔不听话”这个不确定因素控制住了。没有校验的工作流就像没有刹车的车跑得越快越危险。4.3 输出落盘格式统一与目录组织输出环节我遵循一个原则输出格式由工作流决定不由模型决定。模型只负责生成内容格式由执行器统一套模板。这样做的好处是无论模型输出什么格式最终落盘的文件格式都是一致的。我的file_write执行器支持三种输出模式raw直接写模型输出、markdown套Markdown模板、json结构化输出。对于摘要分类这种场景我用markdown模式模板在编排配置里定义- id: save_result type: file_write params: path: {{output_dir}}/{{extract_text.filename}}.md mode: markdown template: | # {{extract_text.filename}} ## 摘要 {{summarize.output}} ## 难度等级 {{classify.output}} ## 处理时间 {{context.timestamp}}目录组织我按“日期/类别/文件名”三层来分。日期用处理当天的日期类别用分类步骤的输出文件名保持原文件名。这样跑完之后输出目录本身就是一份整理好的资料库不需要再手动归类。提示输出目录一定要用绝对路径或者基于项目根目录的相对路径。我一开始用了相对于当前工作目录的路径结果在不同地方执行脚本时输出到了不同位置找文件找了半天。后来统一改成基于配置文件所在目录的绝对路径这个问题就再也没出现过。5. 跑起来之后才发现的坑调度、并发与错误处理5.1 串行够用但什么时候需要并发我一开始是纯串行执行一个文件处理完再处理下一个。对于几十个文件来说串行完全够用因为瓶颈在模型推理速度上不在调度上。但当我需要处理几百个文件时串行就太慢了。并发的实现我走了一条保守路线先做文件级并发不做步骤级并发。也就是说多个文件同时处理但每个文件内部的步骤仍然是串行的。这样做的好处是逻辑简单不需要处理步骤之间的数据竞争问题。实现上用一个线程池每个线程跑一个完整的文件处理流程。from concurrent.futures import ThreadPoolExecutor def process_all_files(file_list, workflow_config, max_workers4): with ThreadPoolExecutor(max_workersmax_workers) as executor: futures { executor.submit(run_workflow, f, workflow_config): f for f in file_list } for future in futures: try: future.result() except Exception as e: print(f处理失败: {futures[future]}, 错误: {e})并发数我设的是4。为什么不是更高因为本地模型推理本身就会占满GPU并发太高反而会导致显存不足或者推理速度下降。4个并发在我的机器上刚好能把GPU利用率拉满再高就没有明显收益了。5.2 错误处理让流程“断而不崩”工作流跑批量任务时最怕的就是一个文件出错导致整个流程崩掉。我的处理策略是步骤级捕获文件级隔离。每个步骤执行时都用try-except包住出错就把错误信息记到日志里然后跳过这个文件继续处理下一个。错误日志我记三个东西出错的文件名、出错的步骤ID、具体的错误信息。有了这三样排查问题基本就是看一眼日志的事。我还会在流程结束时输出一个汇总总共处理了多少个文件成功多少失败多少失败的文件列表。这样跑完一眼就能看出有没有问题。def run_workflow(file_path, config): context {input_path: file_path, timestamp: now()} for step in config[steps]: try: result execute_step(step, context) context[step[id]] result except Exception as e: log_error(file_path, step[id], str(e)) return {status: failed, step: step[id], error: str(e)} return {status: success}还有一个细节中间结果要不要保存。我一开始不保存出错就从头重跑。后来发现有些步骤比如模型调用比较慢重跑成本高。于是加了一个可选的checkpoint机制每个步骤执行完把上下文序列化到临时文件重跑时如果发现有checkpoint就直接加载跳过已完成的步骤。这个机制在调试阶段特别有用改一个后面的步骤不需要把前面的模型调用全部重跑一遍。5.3 日志与可观测性出了问题怎么快速定位日志我分了三个级别INFO记录每个步骤的开始和结束WARN记录重试和降级ERROR记录失败。日志格式是时间戳 | 级别 | 文件名 | 步骤ID | 信息用竖线分隔方便用grep过滤。2025-01-15 10:23:01 | INFO | doc_001.md | extract_text | 开始执行 2025-01-15 10:23:01 | INFO | doc_001.md | extract_text | 完成输出长度 3421 2025-01-15 10:23:02 | INFO | doc_001.md | summarize | 开始执行 2025-01-15 10:23:08 | WARN | doc_001.md | summarize | 首次调用超时重试中 2025-01-15 10:23:15 | INFO | doc_001.md | summarize | 完成输出长度 187有了这个日志排查问题就是一条命令的事grep ERROR workflow.log看有哪些错误grep doc_001 workflow.log看某个文件的完整处理链路。我还会在流程结束时统计每个步骤的平均耗时这样能看出哪个步骤是瓶颈后续优化就有方向了。6. 这套工作流真正省下的到底是什么回到标题里的“六毛钱”。那六毛钱本身其实不是重点重点是它代表的那种按次付费、不可控、不可定制的工具使用方式。我搭这套工作流花了大概两个周末的时间之后每次使用都是零工具成本。但比省钱更重要的是另外三件事。第一是可控性。每个步骤用什么模型、提示词怎么写、输出什么格式、出错怎么处理全部由我决定。我不需要等平台更新功能也不需要因为平台改版而重新适应。这种掌控感是用现成工具永远得不到的。第二是可组合性。这套工作流里的每个执行器都是独立的我可以把llm_call执行器单独拿出来接进别的项目也可以把整个工作流当成一个步骤嵌进更大的流程里。这种灵活性是平台级工具给不了的因为平台级工具的设计目标就是让你留在它的生态里。第三是可学习性。搭这套工作流的过程逼着我把“AI到底是怎么被调用的”“提示词为什么会影响输出”“工作流编排的本质是什么”这些问题想清楚了。这些认知比省下的钱值钱得多。如果你也想搭一套类似的我的建议是从最简单的开始。不要一上来就搞复杂的并发和错误处理先写一个能跑通“读文件-调模型-写文件”三步的脚本跑通了再逐步加功能。我见过太多人一开始就想搭一个“万能工作流”结果卡在架构设计上最后什么都没跑起来。先跑通再优化这是我在实际项目里反复验证过的路径。最后分享一个我常用的调试技巧用固定输入测执行器用变化输入测编排。执行器的测试用同一段文本反复跑看输出是否稳定编排的测试用不同文件跑看流程是否能正确处理各种边界情况。分开测问题定位会快很多。
返回列表