从脚本到系统:AI Agent工作流工程化实战与架构演进

从脚本到系统:AI Agent工作流工程化实战与架构演进

1. 项目概述:从脚本到系统的思维跃迁

最近和不少同行交流,发现一个挺普遍的现象:大家用大模型API写个脚本,调用几个工具,跑起来能完成一个任务,就觉得自己已经搞定了“AI Agent”。这其实是个挺大的误区。一个能跑通的Python脚本,和一个真正可维护、可扩展、能稳定运行的自动化系统,中间隔着一道巨大的鸿沟。我最初也是从写线性脚本开始的,一个文件几百行代码,所有逻辑都揉在一起,今天加个功能,明天改个参数,很快就变成了一团乱麻,调试起来痛苦不堪,更别提交给别人维护了。这个项目,就是记录了我如何把一个简单的、线性的“AI任务执行脚本”,一步步重构、设计,最终演变成一个模块清晰、职责分明、具备一定容错和监控能力的“AI Agent工作流系统”的完整实战过程。这不仅仅是代码层面的优化,更是一种工程化思维的转变。如果你也正困在脚本的泥潭里,或者想系统地了解如何构建一个健壮的AI自动化应用,那么接下来的内容应该能给你带来不少启发。

2. 核心架构演进:从面条代码到分层设计

2.1 线性脚本的典型困境

我们从一个最简单的场景开始:一个自动分析行业报告并生成摘要的脚本。最初的版本可能长这样:

import openai import requests from bs4 import BeautifulSoup def main(): # 1. 读取配置 api_key = “your_key” report_url = “http://example.com/report” # 2. 爬取报告内容 response = requests.get(report_url) soup = BeautifulSoup(response.text, ‘html.parser’) content = soup.get_text() # 3. 调用大模型总结 client = openai.OpenAI(api_key=api_key) prompt = f“请总结以下内容:{content}” response = client.chat.completions.create( model=“gpt-4”, messages=[{“role”: “user”, “content”: prompt}] ) summary = response.choices[0].message.content # 4. 保存结果 with open(‘summary.txt’, ‘w’) as f: f.write(summary) if __name__ == ‘__main__’: main()

这个脚本能工作,但它暴露了所有典型问题:配置硬编码逻辑耦合没有错误处理无法复用。如果你想换个模型、增加一个数据清洗步骤、或者处理多个报告,就必须直接修改这个“主函数”,风险极高。

2.2 迈向模块化的第一步:功能解耦

第一步重构,我们把不同的职责拆分成独立的函数或类。

class ReportFetcher: def fetch(self, url): # 包含重试、超时、异常处理 pass class ContentParser: def parse(self, html): # 包含清洗、提取正文 pass class LLMService: def __init__(self, model, api_key): self.client = openai.OpenAI(api_key=api_key) self.model = model def summarize(self, content): # 包含prompt模板、调用、解析 pass class ResultSaver: def save(self, content, path): # 包含文件操作、格式处理 pass

这样,每个类只负责一件事。但此时,它们之间的调用关系依然是硬编码在另一个“协调器”函数里,流程是固定的。这解决了代码组织问题,但还没解决流程灵活性问题。

2.3 引入工作流引擎思维

当任务变得复杂,例如“获取报告 -> 清洗内容 -> 调用模型A提取关键点 -> 调用模型B生成摘要 -> 存入数据库 -> 发送通知”,固定的代码顺序就很难维护了。这时,我们需要一个“工作流引擎”的概念。它不关心每个步骤具体怎么做,只关心步骤之间的依赖关系执行顺序

我调研了n8n、temporal等,但对于AI Agent场景,它们有时显得过重。我选择借鉴其思想,实现一个轻量级的、基于DAG(有向无环图)的工作流执行器。核心是定义“节点”(Node)和“边”(Edge)。

  • 节点:对应一个具体的处理单元,如FetchReportNode,SummarizeWithLLMNode。每个节点有输入、输出和execute方法。
  • :定义了节点之间数据的流动关系,例如节点A的输出summary_text作为节点B的输入raw_content

这样,整个流程就可以用配置文件(如YAML)或代码来定义,而不是写死在主程序里。

workflow: name: “report_analysis” nodes: - id: fetch type: ReportFetcher params: url: “{input_url}” - id: parse type: ContentParser depends_on: [“fetch”] - id: summarize type: LLMService model: “gpt-4” depends_on: [“parse”] params: prompt_template: “总结:{content}” output: summarize.result

这个架构的转变,是系统可维护性的基石。它允许我们:

  1. 热更新流程:修改YAML文件即可调整业务流程,无需重启服务或修改代码。
  2. 节点复用:同一个LLMService节点可以被不同的工作流使用。
  3. 可视化编排:理论上,可以基于此结构开发一个简单的可视化编辑器(类似Dify工作流、n8n那样)。

注意:在项目初期,不要过度设计。如果你的业务逻辑非常简单且稳定,一个清晰的模块化脚本就足够了。引入工作流引擎会带来额外的复杂度。我的经验是,当手动编排的“if-else”或顺序调用超过3层,或者经常需要调整步骤时,就该考虑工作流了。

3. 核心组件深度解析与实现

3.1 Agent Skill:可复用的能力单元

在AI Agent语境下,工作流中的节点可以进一步抽象为“Skill”(技能)。一个Skill是一个自包含的、可独立测试的能力单元。例如,“网页搜索技能”、“数据库查询技能”、“发送邮件技能”。

我定义的Skill接口通常包含以下部分:

from abc import ABC, abstractmethod from typing import Any, Dict class BaseSkill(ABC): “”“技能基类”“” name: str description: str @abstractmethod def execute(self, input_data: Dict[str, Any], context: Dict[str, Any]) -> Dict[str, Any]: “”“执行技能,返回结果字典”“” pass def get_input_schema(self) -> Dict: “”“声明所需的输入参数格式,用于动态UI生成和验证”“” return {} def get_output_schema(self) -> Dict: “”“声明输出数据的格式”“” return {}

为什么需要Schema?这对于构建自动化系统和低代码平台至关重要。有了输入输出Schema,上游节点可以知道该传递什么数据给这个Skill,下游节点可以知道能接收到什么数据。更重要的是,它可以用于自动生成前端表单,让非开发者也能配置这个Skill的参数。

一个具体的LLM调用Skill实现示例:

class LLMSummarizeSkill(BaseSkill): name = “llm_summarize” description = “使用大模型对文本进行摘要总结” def __init__(self, llm_client, default_model=“gpt-3.5-turbo”): self.client = llm_client self.default_model = default_model def get_input_schema(self): return { “type”: “object”, “properties”: { “text”: {“type”: “string”, “description”: “需要总结的原始文本”}, “model”: {“type”: “string”, “description”: “使用的模型”, “default”: self.default_model}, “max_length”: {“type”: “integer”, “description”: “摘要最大长度”, “default”: 500} }, “required”: [“text”] } def execute(self, input_data, context): text = input_data.get(“text”) model = input_data.get(“model”, self.default_model) max_length = input_data.get(“max_length”, 500) # 构建更鲁棒的prompt prompt = f“””请为以下文本生成一个简洁的摘要,摘要长度不超过{max_length}字。 文本内容: {text} 摘要:”“” try: response = self.client.chat.completions.create( model=model, messages=[{“role”: “user”, “content”: prompt}], temperature=0.2 # 降低随机性,使摘要更稳定 ) summary = response.choices[0].message.content.strip() return {“success”: True, “summary”: summary, “model_used”: model} except Exception as e: # 记录错误上下文,便于排查 return {“success”: False, “error”: str(e), “input_snapshot”: input_data}

这个实现相比简单脚本,增加了参数验证结构化返回(包含成功状态和元数据)、异常捕获。返回的字典会成为下一个节点的输入。

3.2 工作流引擎:轻量级DAG执行器

有了Skill,我们需要一个引擎来按顺序或并行地执行它们。下面是一个极度简化的核心执行逻辑,演示如何解析依赖并执行。

class WorkflowEngine: def __init__(self): self.skills = {} # 技能注册表 self.graph = {} # 邻接表表示的DAG def register_skill(self, skill: BaseSkill): self.skills[skill.name] = skill def build_graph(self, workflow_def: Dict): “”“根据定义构建执行图”“” # workflow_def 包含节点列表和依赖关系 for node in workflow_def[“nodes”]: self.graph[node[“id”]] = { “skill”: node[“type”], “params”: node.get(“params”, {}), “depends_on”: node.get(“depends_on”, []), “output”: None } def execute(self, initial_input: Dict) -> Dict: “”“执行工作流”“” # 1. 拓扑排序,确定执行顺序 execution_order = self._topological_sort() context = {“initial_input”: initial_input} node_outputs = {} # 2. 按顺序执行每个节点 for node_id in execution_order: node_info = self.graph[node_id] skill_name = node_info[“skill”] skill = self.skills.get(skill_name) if not skill: raise ValueError(f“Skill {skill_name} not registered”) # 准备输入:合并初始输入、参数模板、上游节点输出 input_data = self._prepare_input(node_info, node_outputs, context) # 执行技能 result = skill.execute(input_data, context) node_outputs[node_id] = result # 如果某个节点失败,可以在这里定义重试或熔断策略 if not result.get(“success”, True): # 例如:重试3次,或触发备用流程 self._handle_node_failure(node_id, result, context) # 3. 收集最终输出 return self._collect_outputs(node_outputs, workflow_def)

关键点解析

  • 拓扑排序:确保依赖项先执行。这是DAG执行的核心。
  • 输入准备 (_prepare_input):这是最容易出问题的地方。需要处理参数替换,例如将“url”: “{input_report_url}”中的{input_report_url}替换为实际值。我实现了一个简单的模板渲染器,支持从initial_inputnode_outputscontext中取值。
  • 错误处理 (_handle_node_failure):简单的脚本遇到错误就崩溃。系统则需要策略。我实现了几个级别:1)节点级重试(针对网络抖动);2)备用节点切换(如主LLM服务失败,换备用API);3)工作流级回退(记录错误,发送告警,保存中间状态供人工干预)。

3.3 状态管理与上下文传递

线性脚本的变量是全局或函数内传递的。在工作流中,数据在节点间流动,需要一个统一的“上下文”(Context)来管理。我的Context对象通常包含:

  • Execution ID:本次工作流执行的唯一标识,用于串联所有日志和监控。
  • Input/Output Data:每个节点的输入输出快照。
  • Skill Registry:可用的技能引用。
  • Configuration:全局配置,如API密钥(通过环境变量或配置中心注入,绝不硬编码)。
  • Logger:附带了Execution ID的日志器,方便追踪。

上下文对象作为参数传递给每个Skill的execute方法,Skill可以将一些中间状态或全局信息存入其中,供后续节点读取。这避免了使用全局变量,也让数据流更清晰。

4. 可维护性提升的关键实践

4.1 配置外部化与安全管理

所有可变的配置项都必须从代码中剥离。我使用分层配置:

  1. 环境变量:用于敏感信息(API Keys,数据库密码)。使用python-dotenv加载。
  2. 配置文件(YAML/JSON):用于工作流定义、模型参数、超时时间等。
  3. 配置中心(可选):在微服务架构下,使用Consul或Nacos实现动态配置更新。

对于AI API密钥,我强烈建议使用密钥管理服务(如AWS KMS, HashiCorp Vault),或者在应用中实现一个简单的密钥轮换代理。绝对不要在代码或配置文件中提交明文密钥。

4.2 日志、监控与可观测性

脚本可以print,系统必须有完善的日志。我采用结构化日志(JSON格式),每个日志条目都包含execution_id,node_id,skill_name,level,timestamp,message。这样可以通过execution_id轻松过滤出一次完整工作流的所有日志。

监控方面,至少需要:

  • 性能指标:每个Skill的执行耗时、成功率。使用Prometheus客户端库暴露指标,用Grafana展示。
  • 业务指标:如“每日处理报告数”、“摘要生成平均长度”。
  • 告警:对失败率升高、耗时异常进行告警(集成到钉钉、Slack或PagerDuty)。

4.3 测试策略:从单元到集成

可维护性的核心是可信赖的测试。

  • 单元测试:针对每个Skill。Mock掉所有外部依赖(网络请求、数据库、LLM API)。测试其内部逻辑、输入验证和错误处理。
    def test_summarize_skill_with_mock(): mock_client = Mock() mock_client.chat.completions.create.return_value.choices[0].message.content = “这是一个测试摘要” skill = LLMSummarizeSkill(mock_client) result = skill.execute({“text”: “测试文本”}, {}) assert result[“success”] is True assert “测试摘要” in result[“summary”]
  • 集成测试:测试两个或多个Skill的串联,使用测试专用的外部服务(如本地Mock Server模拟LLM API)。
  • 工作流测试:将整个工作流定义作为测试用例,使用固定的输入,断言最终的输出是否符合预期。这能有效防止在修改某个节点后,意外破坏整个流程。

4.4 版本控制与CI/CD

工作流定义文件(YAML)和Skill代码一起纳入Git版本控制。每次修改工作流,都应有对应的Commit信息说明变更原因。CI/CD流水线应自动运行全套测试,并在测试通过后,将工作流定义部署到相应环境(开发、测试、生产)。对于复杂的系统,甚至可以做到“蓝绿部署”工作流版本,逐步切流。

5. 实战案例:构建一个智能内容处理流水线

假设我们要构建一个系统:自动从预设的RSS源抓取科技文章,进行内容摘要,提取关键词,判断情感倾向,最后将结果推送到Notion知识库。

5.1 工作流定义

name: “tech_digest_pipeline” version: “1.0” nodes: - id: fetch_feeds type: rss_fetcher params: feed_urls: - “https://example.com/feed1” - “https://example.com/feed2” schedule: “0 9 * * *” # 每天上午9点执行 - id: filter_articles type: content_filter depends_on: [“fetch_feeds”] params: min_length: 500 keywords: [“AI”, “机器学习”, “开源”] - id: summarize type: llm_summarizer depends_on: [“filter_articles”] params: model: “gpt-4-turbo-preview” style: “bullet_points” - id: extract_tags type: keyword_extractor # 可以是用LLM,也可以是传统NLP库 depends_on: [“filter_articles”] - id: sentiment type: sentiment_analysis depends_on: [“filter_articles”] - id: format_output type: notion_formatter depends_on: [“summarize”, “extract_tags”, “sentiment”] params: database_id: “YOUR_DATABASE_ID” - id: upload type: notion_poster depends_on: [“format_output”]

这个定义清晰地描述了7个节点的依赖关系。filter_articles依赖于fetch_feedssummarizeextract_tagssentiment并行依赖于filter_articles,最后format_outputupload串行执行。

5.2 关键Skill实现细节

1. RSS Fetcher Skill:

  • 使用feedparser库。
  • 实现重试机制:网络请求失败时,使用指数退避策略重试。
  • 增量抓取:记录已抓取文章的ID或发布时间,避免重复处理。
  • 超时控制:为每个feed源设置独立超时,防止一个慢源阻塞整个流程。

2. LLM Summarizer Skill:

  • Prompt工程:设计稳定的Prompt模板,将指令、示例、待总结文本清晰分隔。
  • Token管理:计算输入文本的Token数,如果超过模型上限,自动采用“分块总结再合并”的策略。
  • Fallback策略:如果主模型(如GPT-4)调用失败或超时,自动降级到备用模型(如Claude Haiku或本地模型)。

3. Notion Poster Skill:

  • 使用Notion官方SDK。
  • 处理速率限制:Notion API有速率限制,需要在Skill中加入延迟和队列机制。
  • 错误恢复:如果创建页面失败,记录失败的数据和原因,以便后续手动或自动重试。

5.3 部署与调度

这个工作流需要定时触发。我使用了Celery + Redis作为分布式任务队列。

  • 将整个工作流封装成一个Celery Task。
  • 使用Celery Beat根据schedule字段(如crontab)定时发起任务。
  • 执行引擎作为Task的内部逻辑被调用。
  • 这样做的好处是天然获得了队列、重试、监控(Flower)和能力扩展(增加Worker)。

对于更复杂、长期运行的工作流(例如涉及人工审批环节),可以考虑使用TemporalAirflow这类更强大的工作流引擎。它们提供了持久化执行状态、异步等待、版本管理等功能。

6. 常见问题与排查指南

在实际开发和运维中,我遇到了无数坑。这里总结几个最典型的:

问题1:节点执行顺序混乱或依赖错误。

  • 现象:节点B需要节点A的输出,但实际执行时A还没跑完。
  • 排查
    1. 检查工作流定义中的depends_on字段是否正确。
    2. 在引擎的_topological_sort方法后打印执行顺序,验证是否符合预期。
    3. 检查是否有循环依赖(DAG不能有环),可以在建图后运行一次环检测算法。
  • 解决:确保依赖关系声明正确,并使用可靠的拓扑排序算法。

问题2:上下文数据传递丢失或格式错误。

  • 现象:节点C收不到节点B传递的数据,或收到None
  • 排查
    1. _prepare_input方法中增加详细日志,打印出为每个节点准备输入数据时的完整字典。
    2. 检查上游节点的输出Schema是否与下游节点的输入Schema匹配。特别是键名是否一致。
    3. 检查模板字符串的替换逻辑,{variable}中的variable是否能在上下文中找到。
  • 解决:实现一个严格的Schema验证层,在节点执行前验证输入数据格式。为模板渲染失败添加明确的错误信息。

问题3:外部API调用不稳定导致整个工作流失败。

  • 现象:调用OpenAI API或爬取网站时随机超时或失败。
  • 解决策略(层层递进)
    1. 重试:在Skill内部对瞬态错误(网络超时、5xx错误)进行有限次重试(如3次),并加入随机延迟(jitter)。
    2. 熔断:如果某个外部服务连续失败多次,暂时“熔断”对该服务的调用,快速失败并执行降级逻辑,过一段时间再尝试恢复。
    3. 降级:准备备用方案。例如,主LLM服务失败,则调用另一个LLM API;或者使用一个更简单的本地文本摘要算法。
    4. 异步与超时:为所有外部调用设置合理的超时时间,并使用异步IO防止阻塞。

问题4:工作流执行时间长,难以调试。

  • 现象:一个工作流跑几个小时,中间出错很难定位。
  • 解决
    1. 持久化中间状态:在每个节点执行后,将其输入、输出、状态(成功/失败)连同execution_id一起存入数据库(如MongoDB)或对象存储。这样可以通过execution_id随时复盘整个执行过程。
    2. 实现检查点(Checkpoint):对于耗时极长的工作流,可以在关键节点后设置检查点。如果系统崩溃重启,可以从最后一个成功的检查点恢复,而不是从头开始。
    3. 可视化追踪:开发一个简单的管理后台,通过execution_id查询并图形化展示工作流的执行路径、每个节点的状态和耗时,一目了然。

问题5:Prompt效果不稳定,输出格式不符合下游节点要求。

  • 现象:LLM节点有时返回JSON,有时返回纯文本,导致解析失败。
  • 解决
    1. 结构化输出:强制要求LLM以指定格式(如JSON)输出,并在Prompt中给出清晰的示例。使用OpenAI的response_format参数(如果模型支持)。
    2. 输出后处理:在Skill的execute方法中,增加一个post_process步骤,用于清洗和格式化LLM的返回结果。例如,用正则表达式从文本中提取JSON,或者对非标准格式进行容错解析。
    3. Prompt版本化:将Prompt模板也作为配置管理起来,修改Prompt像修改代码一样需要经过评审和测试。

从线性脚本到自动化系统的旅程,本质上是从“让代码跑起来”到“让系统可靠、可维护、可演进”的思维升级。这个过程里,工具和框架的选择固然重要,但更关键的是对边界职责数据流的清晰定义。我的体会是,不要试图一开始就设计一个完美无缺的庞大系统,而是从最痛的“脚本维护之痛”出发,一步步抽象、解耦、加固。先让一个核心流程以模块化的方式稳定跑起来,然后自然地,你会发现哪些地方需要工作流引擎,哪些地方需要状态管理,哪些监控指标必不可少。最终,你会收获的不仅是一个好用的AI Agent系统,更是一套应对复杂性的工程方法。