ARTICLE DETAIL

资讯详情

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

AI工作流四层架构实战:从接入到落地的稳定工程化指南

AI工作流四层架构实战:从接入到落地的稳定工程化指南 1. 为什么我把整套AI工作流拆成了四层1.1 从一次“跑不通”的深夜调试说起去年冬天我接了个私活帮一个做电商的朋友搭一套自动化的商品文案生成系统。需求听起来不复杂每天定时抓取店铺新品信息调用大模型生成卖点文案审核后自动发布到后台。我当时想这不就是“爬虫 调API 写数据库”三件套吗两天搞定。结果我整整折腾了十一天。问题出在哪不是某个单点技术不会而是链路一长任何一环的抖动都会把整条流水线拖垮。爬虫被限流了重试逻辑没写好任务卡死模型接口偶尔超时没有降级方案整个批次全挂生成的内容格式不稳定JSON解析报错下游写库直接崩。最要命的是我根本不知道是哪一步出的问题日志散落在三个脚本里排查一次要半小时。那次之后我彻底想明白一件事AI工作流的核心难点从来不是“能不能调通模型”而是“能不能稳定、可观测、可恢复地跑完一整条链路”。这也是我后来把整套东西拆成四层的起点。1.2 四层架构接入层、编排层、推理层、落地层我现在做任何AI工作流脑子里都先画这四层。你可以把它理解成一条工厂流水线接入层负责“原料进厂”。数据从哪来是定时抓取、消息队列推送还是用户手动触发这一层要解决的是数据获取和初步清洗。编排层负责“工序调度”。哪个任务先跑、哪个后跑、失败了怎么办、并发多少、超时多久。这一层是整个工作流的大脑。推理层负责“核心加工”。调用大模型、做分类、做抽取、做生成。这一层要处理的是模型的不确定性。落地层负责“成品出厂”。结果写库、发通知、存文件、推送到业务系统。为什么这么分因为每一层的失败模式完全不同混在一起就没法针对性治理。接入层怕限流和脏数据编排层怕死锁和雪崩推理层怕超时和幻觉落地层怕格式错和重复写。分开之后每一层我都能单独做重试、单独做监控、单独做降级。提示如果你现在的工作流还是“一个Python脚本从头跑到尾”强烈建议先做这一步拆分。不用改代码先在纸上把每个环节归到四层里你会立刻发现哪些地方是裸奔的。1.3 编排层选型为什么我最后没上Airflow说到编排很多人第一反应是Airflow、Prefect、Dagster这些专业工具。我试过Airflow本地Docker起了一套DAG写起来确实优雅。但问题是我这个场景是轻量、高频、单机为主的Airflow的调度器、元数据库、Web UI一整套下来资源占用比我业务本身还大。后来我选了一个更朴素的方案用Python的concurrent.futures 一个轻量的任务状态表。核心思路是每个任务是一个函数函数执行前后往SQLite里写状态主调度器轮询状态表决定下一步跑什么。听起来很土但它有几个实打实的好处依赖极少pip install两三个包就能跑状态全在本地文件里出问题直接打开看想加并发就加线程池想加超时就用future.result(timeout...)当然如果你的任务量到了每天几十万条、需要跨机器调度那还是老老实实上专业工具。选型的核心不是“哪个最先进”而是“哪个的复杂度匹配我当前的规模”。我见过太多人为了一个每天跑几百条的任务上了K8s最后维护成本比业务代码还高。1.4 推理层的核心矛盾稳定性和灵活性的拉扯推理层是最容易让人上头的地方。一开始我总想把prompt调得完美让模型每次都输出标准JSON。后来发现这是徒劳的——大模型的输出天然带有不确定性你越是想用prompt锁死它它越会在边缘case上给你惊喜。我现在的做法是“宽进严出”prompt里给足示例和格式要求但代码层面必须做输出校验和修复。比如要求输出JSON那我就写一个容错解析器遇到多余的文字就正则提取遇到单引号就替换成双引号遇到截断就补全括号。实测下来这套组合拳能把解析成功率从70%拉到99%以上。这个思路其实和做前端表单校验一样你不能指望用户永远输入正确但你可以让系统在用户输错时优雅地兜住。2. 接入层数据抓取与清洗的实操细节2.1 用Python爬虫做数据接入的三个关键参数接入层我大部分时候用requestsBeautifulSoup偶尔上playwright处理动态页面。这里不展开讲爬虫怎么写重点说三个新手最容易忽略但老手必调的参数。第一个是超时时间。requests.get(url, timeout10)这个10秒不是随便写的。我的经验是连接超时设5秒读取超时设15秒分开设置。因为连接失败通常是网络问题快速失败快速重试读取慢可能是对方服务器在慢慢吐数据给足时间。你可以这样写import requests from requests.adapters import HTTPAdapter from urllib3.util.retry import Retry session requests.Session() retry Retry(total3, backoff_factor1, status_forcelist[429, 500, 502, 503]) session.mount(https://, HTTPAdapter(max_retriesretry)) resp session.get(url, timeout(5, 15))第二个是请求间隔。很多人被限流就是因为请求太密。我的做法是每次请求后sleep(random.uniform(1, 3))加随机抖动。别小看这个随机固定间隔反而容易被识别出是机器行为。第三个是User-Agent轮换。准备一个列表每次请求随机选一个。这不是为了做什么见不得光的事纯粹是为了让请求看起来更自然降低被误伤的概率。2.2 数据清洗把脏数据挡在推理层之外接入层拿到的原始数据十有八九是脏的。HTML标签、多余空格、编码错误、字段缺失什么都有。我的原则是能在接入层清洗的绝不留给推理层。因为推理层调用模型是要花钱的你把一堆垃圾喂给模型既浪费token又影响输出质量。清洗我一般分三步走结构化提取用CSS选择器或XPath把需要的字段抠出来不要整页丢给模型。文本规范化去掉HTML标签、合并连续空白、统一全半角、处理特殊字符。字段校验检查必填字段是否为空长度是否超限格式是否符合预期。这里有个小技巧把清洗规则写成配置而不是硬编码在代码里。比如哪些字段必填、最大长度多少、用什么正则校验全部放到一个YAML文件里。这样业务方改需求时你改配置就行不用动代码重新部署。fields: title: required: true max_length: 200 price: required: true pattern: ^\d(\.\d{1,2})?$ description: required: false max_length: 20002.3 接入层的容错设计断点续传和去重接入层最怕的是跑到一半挂了重启后从头再来。我的做法是每处理完一条数据就记录一个游标存在SQLite或Redis里。重启时从游标位置继续而不是从头开始。去重也很关键。同一个商品可能被重复抓取如果不去重推理层就会重复调用模型白白烧钱。我的去重策略是对业务主键做哈希存到一个集合里处理前先查集合存在就跳过。简单粗暴但有效。注意去重集合如果放在内存里重启就丢了。建议持久化到本地文件或数据库虽然慢一点但省心。3. 编排层任务调度与状态管理的落地方法3.1 任务状态机让每个任务都有迹可循编排层的核心是状态管理。我给每个任务定义了五个状态pending待执行、running执行中、success成功、failed失败、retrying重试中。每次状态变更都写一条记录到数据库包含任务ID、状态、时间戳、错误信息。这样做的好处是任何时候我打开数据库就能看到整条链路的实时状态。哪个任务卡在running超过预期时间哪个任务failed了三次一目了然。状态流转的逻辑我用一个简单的函数封装def update_task_status(task_id, status, errorNone): conn sqlite3.connect(tasks.db) conn.execute( UPDATE tasks SET status?, error?, updated_at? WHERE id?, (status, error, datetime.now(), task_id) ) conn.commit() conn.close()别嫌它土能跑通、能排查、能恢复就是好编排。3.2 重试策略指数退避加最大次数重试不是简单地“失败了再跑一次”。我见过有人写while True: try: ... except: pass这是灾难。正确的重试要满足三个条件有最大次数、有退避间隔、有失败兜底。我的标准配置是最大重试3次间隔用指数退避第一次等2秒第二次等4秒第三次等8秒。这样既能应对临时抖动又不会在对方服务彻底挂掉时疯狂重试把人家打垮。import time def retry_with_backoff(func, max_retries3, base_delay2): for attempt in range(max_retries): try: return func() except Exception as e: if attempt max_retries - 1: raise delay base_delay * (2 ** attempt) time.sleep(delay)如果三次都失败任务标记为failed进入人工介入队列。不要无限重试那只会掩盖问题。3.3 并发控制线程池大小怎么定并发能提速但并发太高会把下游打挂。我的经验值是IO密集型任务如调API线程数设为CPU核数的5到10倍CPU密集型任务线程数不超过CPU核数。假设你的机器是4核调模型API属于IO密集型那线程池设20到40比较合适。但这不是绝对的还要看下游能承受多少QPS。我一般先设一个保守值比如10然后压测逐步往上加观察错误率和响应时间找到拐点。from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor(max_workers10) as executor: futures [executor.submit(process_item, item) for item in items] for future in futures: try: result future.result(timeout60) except TimeoutError: update_task_status(item.id, failed, timeout)提示线程池的max_workers不是越大越好。我踩过的坑是设了50结果下游API直接返回429整批任务全挂。后来降到15反而跑得更稳。3.4 超时与熔断给每个任务装上保险丝超时是必须的。一个任务卡住不返回会占着线程不放线程池很快就被占满整个系统瘫痪。我给每个任务设了硬超时比如调模型API设60秒爬虫设30秒。超时就标记失败释放线程。熔断是进阶玩法。如果某个下游服务连续失败超过阈值比如10次里失败8次就暂时停止调用它直接走降级逻辑过一段时间再试探性恢复。这个用pybreaker库几行代码就能实现。熔断的意义在于当下游已经挂了你不要再往上撞把资源留给还能跑的任务。4. 推理层让大模型稳定输出的工程技巧4.1 Prompt工程结构化输入比华丽措辞更重要很多人写prompt喜欢堆形容词“请你作为一个专业的、资深的、经验丰富的专家……”。实测下来这些修饰词对输出质量的提升微乎其微真正有用的是结构化的输入和明确的输出格式要求。我现在的prompt模板长这样任务根据商品信息生成卖点文案 输入 - 商品名称{title} - 价格{price} - 核心参数{specs} 要求 1. 输出JSON格式包含字段headline标题、points卖点数组3条、summary总结 2. 每条卖点不超过20字 3. 不要编造输入中没有的信息 示例输出 {headline: ..., points: [..., ..., ...], summary: ...}关键点在于给示例、给格式、给约束。示例让模型知道你要什么风格格式让解析变得简单约束不要编造能减少幻觉。4.2 输出解析容错解析器的写法即使你要求了JSON模型也可能给你返回带markdown代码块的、带解释文字的、或者格式微错的JSON。所以解析器必须容错。我的解析流程是先尝试直接json.loads失败则用正则提取{...}之间的内容再解析再失败则替换常见错误单引号转双引号、去尾逗号后解析还失败就记录原始输出标记为待人工处理import json import re def parse_json_safe(text): try: return json.loads(text) except json.JSONDecodeError: pass match re.search(r\{.*\}, text, re.DOTALL) if match: cleaned match.group() cleaned cleaned.replace(, ) cleaned re.sub(r,\s*}, }, cleaned) try: return json.loads(cleaned) except json.JSONDecodeError: pass return None这套组合拳下来解析成功率能到99%以上。剩下那1%人工兜底。4.3 模型选型不是越贵越好推理层的模型选型我的原则是按任务难度分级。简单的分类、抽取任务用便宜的小模型就够了复杂的生成、推理任务再上大模型。比如商品文案生成我试过几个模型最后发现中等规模的模型在“卖点提炼”这个任务上效果和大模型差距不大但成本只有三分之一。先用小模型跑一批人工评估质量质量达标就用小模型不达标再升级。这个试错成本很低但长期省下来的钱很可观。另外批处理能省钱。如果任务不要求实时返回把多条数据打包成一个请求发给模型比一条一条调便宜得多。很多模型API都支持批量输入值得研究一下。4.4 缓存同样的输入不要调两次模型推理层最容易被忽视的优化是缓存。同一个商品、同样的prompt如果之前已经生成过文案直接读缓存就行没必要再调一次模型。我的缓存键是hash(prompt input)缓存值是模型输出。用Redis或本地文件都行。实测下来在商品文案这种场景缓存命中率能到30%以上直接省掉三分之一的推理成本。注意缓存要设过期时间。商品信息可能更新缓存太旧会导致输出过时。我一般设7天。5. 落地层结果输出与业务对接5.1 写库批量写入比逐条快十倍落地层最常见的问题是写库太慢。一条一条INSERT几千条数据要跑好几分钟。改成批量写入速度能快十倍以上。def batch_insert(conn, table, rows, batch_size500): for i in range(0, len(rows), batch_size): batch rows[i:ibatch_size] placeholders ,.join([?] * len(batch[0])) conn.executemany( fINSERT INTO {table} VALUES ({placeholders}), batch ) conn.commit()批量大小我一般设500。太小了频繁提交开销大太大了单次事务太重容易锁表。500是个比较稳的中间值。5.2 幂等性重复执行不会产生脏数据落地层必须保证幂等。什么意思就是同一个任务跑一次和跑十次结果是一样的。这在工作流里特别重要因为重试是常态。实现幂等最简单的方法是用业务主键做唯一约束。写库时用INSERT OR REPLACE或ON CONFLICT DO UPDATE这样重复写入不会产生重复记录。INSERT INTO products (id, name, copywriting) VALUES (?, ?, ?) ON CONFLICT(id) DO UPDATE SET name excluded.name, copywriting excluded.copywriting;5.3 通知与告警让问题主动找你落地层跑完得让人知道结果。我的做法是分级通知全部成功发一条汇总消息有失败的发详细错误列表失败率超过阈值直接打电话告警。通知渠道看团队习惯邮件、企业微信、钉钉都行。关键是要包含足够的信息任务ID、成功数、失败数、失败原因、日志链接。别只发一句“任务完成”那等于没发。6. 常见问题与排查技巧实录6.1 任务卡死不动了怎么办这是最高频的问题。排查顺序是看状态表哪个任务卡在running看该任务的日志最后一条输出是什么如果是调API卡住检查超时设置是否生效如果是数据库锁检查是否有长事务未提交我遇到过一次任务卡在running两小时。查下来是爬虫请求一个不存在的域名DNS解析一直不返回而我没设连接超时。加上timeout(5, 15)后问题解决。6.2 模型输出格式不稳定怎么破前面说的容错解析器能解决大部分问题。如果还是频繁失败说明prompt需要优化。我的经验是在prompt里加一个“只输出JSON不要任何其他文字”的强约束并且把示例输出放在最后。模型对最后的内容印象最深。6.3 内存越跑越高最后OOM这是典型的资源泄漏。常见原因有两个一是线程池没关闭二是缓存没设上限。线程池用with语句管理缓存用LRU策略限制大小。from functools import lru_cache lru_cache(maxsize1000) def get_cached_result(key): ...6.4 常见问题速查表问题现象可能原因排查方法解决方案任务卡在running超时未设或未生效查日志最后输出加超时参数解析频繁失败prompt格式约束弱看原始输出强化格式要求容错解析内存持续增长资源未释放监控内存曲线用with管理资源LRU缓存写库慢逐条写入看SQL执行时间改批量写入重复数据无幂等设计查主键重复加唯一约束upsert下游限流并发过高看429错误率降并发加退避6.5 我踩过的三个坑第一个坑是日志打太多。一开始我把每条数据的完整内容都打进日志结果日志文件一天几个G磁盘直接满了。后来改成只打关键信息和错误详情正常流程只打计数。第二个坑是重试没有区分错误类型。有些错误重试有用网络抖动有些错误重试没用参数错误。我后来加了一个判断只有超时和5xx错误才重试4xx错误直接失败。第三个坑是没有做端到端测试。有一次改了一个字段名上游改了下游没改跑了一整天才发现数据全错。后来我加了一个小批量的端到端测试每次部署前先跑10条数据验证全链路。7. 环境搭建从零把工作流跑起来7.1 Python环境配置的稳妥方案环境这块我推荐用conda建独立环境不要用系统自带的Python。原因很简单系统Python被各种系统工具依赖你乱装包可能把系统搞崩。conda create -n aiworkflow python3.10 conda activate aiworkflow pip install requests beautifulsoup4 openai sqlite3Python版本选3.10兼容性好新特性也够用。别追最新版很多库还没适配。如果你用VSCode装个Python插件在设置里把默认解释器指向conda环境。PyCharm的话在项目设置里配置解释器路径。这些基础操作网上教程很多不展开。7.2 项目目录结构建议ai_workflow/ ├── config/ │ ├── settings.yaml │ └── prompts/ ├── src/ │ ├── ingest/ │ ├── orchestrate/ │ ├── inference/ │ └── output/ ├── data/ │ ├── tasks.db │ └── cache/ ├── logs/ └── tests/按四层架构分目录每层一个包。配置和代码分离数据和日志独立。这个结构不复杂但能让项目长期可维护。7.3 最小可运行示例最后给一个最小可运行的骨架把四层串起来import sqlite3 from concurrent.futures import ThreadPoolExecutor def ingest(): return [{id: 1, title: 测试商品}] def infer(item): return {id: item[id], copy: 生成的文案} def output(results): conn sqlite3.connect(data/tasks.db) conn.executemany( INSERT OR REPLACE INTO results VALUES (?, ?), [(r[id], r[copy]) for r in results] ) conn.commit() def main(): items ingest() with ThreadPoolExecutor(max_workers10) as executor: results list(executor.map(infer, items)) output(results) if __name__ __main__: main()这个骨架跑通了再往里填重试、状态管理、容错解析就是一套完整的生产级工作流。我在实际项目里最大的体会是AI工作流的难点不在AI在工作流。模型调用本身很简单难的是让整条链路在真实环境的各种意外下还能稳定跑完。把四层拆清楚每层做好自己的容错比追求某个单点的“高级技术”有用得多。
返回列表