# 企业级Agent自动化流水线实战
随着大模型技术的成熟,企业开始将AI Agent从“实验性玩具”转向“生产力工具”。但单点Agent在复杂业务场景中往往表现不佳:上下文丢失、工具调用混乱、缺乏安全管控。企业级Agent自动化流水线正是为了解决这些问题而生——它像一条工业生产线,将多个Agent、工具、审批节点、监控机制串联起来,形成稳定、可审计、可扩展的自动化系统。
行业背景: 2024-2025年,金融、医疗、电商等行业已出现大量Agent流水线案例,例如:
1. 设计多Agent协作架构,能根据业务需求拆分任务并编排流水线 2. 整合5种以上外部工具(数据库、API、文件系统、Webhook等) 3. 配置人工审批与异常回退机制,确保关键节点可控 4. 搭建监控与告警体系,实时掌握流水线运行状态 5. 优化流水线性能,处理高并发场景(如每秒100+请求)
---
原理讲解
流水线本质是一个有向无环图(DAG),每个节点是一个Agent或工具。架构设计核心原则:
架构对比
| 方案 | 单体Agent | 简单链式 | DAG编排(推荐) | |------|-----------|----------|------------------| | 复杂度 | 低 | 中 | 高 | | 扩展性 | 差 | 一般 | 优秀 | | 可观测性 | 差 | 一般 | 优秀 | | 典型工具 | LangChain | LangGraph | Temporal / Airflow + Agent |
代码示例:使用LangGraph定义DAG流水线
from langgraph.graph import StateGraph, END
from typing import TypedDict, List
import json
# 定义状态结构
class PipelineState(TypedDict):
user_query: str
intent: str
search_results: str
draft_response: str
approval_status: str
final_response: str
# 定义节点函数(简化版)
def intent_classifier(state: PipelineState):
# 调用LLM进行意图识别
intent = "complaint" # 假设结果
return {"intent": intent}
def knowledge_search(state: PipelineState):
# 根据意图查询知识库
results = "相关FAQ内容..."
return {"search_results": results}
def draft_generator(state: PipelineState):
draft = f"基于查询'{state['user_query']}',建议回复:{state['search_results']}"
return {"draft_response": draft}
def approval_check(state: PipelineState):
# 模拟审批:高敏感意图需要人工
if state["intent"] == "complaint":
return {"approval_status": "pending_human"}
return {"approval_status": "auto_approved"}
def human_review(state: PipelineState):
# 此处实际会等待Webhook回调
print("等待人工审批...")
return {"final_response": state["draft_response"]}
def final_response(state: PipelineState):
return {"final_response": state["draft_response"]}
# 构建图
builder = StateGraph(PipelineState)
builder.add_node("classifier", intent_classifier)
builder.add_node("search", knowledge_search)
builder.add_node("draft", draft_generator)
builder.add_node("approval", approval_check)
builder.add_node("human", human_review)
builder.add_node("final", final_response)
# 定义边
builder.set_entry_point("classifier")
builder.add_edge("classifier", "search")
builder.add_edge("search", "draft")
builder.add_edge("draft", "approval")
# 条件分支
builder.add_conditional_edges(
"approval",
lambda state: "human" if state["approval_status"] == "pending_human" else "final",
{"human": "human", "final": "final"}
)
builder.add_edge("human", "final")
builder.add_edge("final", END)
# 编译并运行
app = builder.compile()
result = app.invoke({"user_query": "我要投诉产品质量问题"})
print(result["final_response"])
模块2:多工具链整合——统一工具接口
原理讲解
企业级流水线需要调用多种外部系统:数据库、REST API、文件系统、消息队列、搜索引擎等。核心是设计统一的工具抽象层,每个工具封装为函数,输入输出遵循标准格式。
关键技术:工具注册与调用模式
from typing import Dict, Any, Callable
import requests
import sqlite3
# 工具注册表
class ToolRegistry:
def __init__(self):
self._tools: Dict[str, Callable] = {}
def register(self, name: str, func: Callable, description: str = ""):
self._tools[name] = func
print(f"已注册工具: {name} - {description}")
def call(self, name: str, params: Dict[str, Any]) -> Any:
if name not in self._tools:
raise ValueError(f"工具 {name} 未注册")
return self._tools[name](**params)
# 示例工具实现
def query_database(sql: str, db_path: str = "orders.db") -> list:
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
cursor.execute(sql)
results = cursor.fetchall()
conn.close()
return results
def call_external_api(url: str, method: str = "GET", headers: dict = None) -> dict:
resp = requests.request(method, url, headers=headers)
resp.raise_for_status()
return resp.json()
def send_webhook(url: str, payload: dict) -> str:
resp = requests.post(url, json=payload)
return f"Webhook sent, status: {resp.status_code}"
# 注册工具
registry = ToolRegistry()
registry.register("database_query", query_database, "执行SQL查询")
registry.register("api_call", call_external_api, "调用外部REST API")
registry.register("webhook", send_webhook, "发送Webhook通知")
# Agent调用示例
def agent_tool_call(state: dict) -> dict:
# 假设LLM决定调用数据库查询
result = registry.call("database_query", {"sql": "SELECT * FROM orders LIMIT 5"})
state["db_result"] = result
return state
模块3:审批与人工介入机制
原理讲解
关键节点需要人工决策时,Agent不能继续执行。设计模式:
- **暂停点:** Agent执行到特定节点时,将上下文持久化到数据库,等待外部信号(Webhook、轮询)
- **超时处理:** 设置审批超时时间,超时后触发降级策略(如自动拒绝或转交)
- **审计日志:** 所有审批操作记录到不可篡改的日志中
代码示例:基于Redis的暂停/恢复机制
import redis
import json
from datetime import datetime
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
def pause_pipeline(pipeline_id: str, state: dict, pending_node: str):
"""将流水线暂停,等待人工审批"""
state["_pending_node"] = pending_node
state["_paused_at"] = datetime.now().isoformat()
r.set(f"pipeline:{pipeline_id}:state", json.dumps(state))
r.set(f"pipeline:{pipeline_id}:status", "pending_approval")
# 发送通知给审批人
send_webhook("https://approval-system/webhook", {
"pipeline_id": pipeline_id,
"action": "approval_required",
"context": state.get("draft_response", "")
})
def resume_pipeline(pipeline_id: str, decision: str, reviewer_note: str = ""):
"""人工审批后恢复流水线"""
state_str = r.get(f"pipeline:{pipeline_id}:state")
if not state_str:
raise ValueError("流水线状态不存在")
state = json.loads(state_str)
state["_review_decision"] = decision
state["_reviewer_note"] = reviewer_note
state["_reviewed_at"] = datetime.now().isoformat()
# 根据决策更新状态
if decision == "approve":
state["approval_status"] = "approved"
else:
state["approval_status"] = "rejected"
r.set(f"pipeline:{pipeline_id}:state", json.dumps(state))
r.set(f"pipeline:{pipeline_id}:status", "resumed")
return state
模块4:监控与可观测性
原理讲解
流水线运行需要三方面监控:
- **Metrics(指标):** 每个节点耗时、成功率、并发数(Prometheus + Grafana)
- **Tracing(链路):** 跟踪一次请求经过的所有节点(OpenTelemetry)
- **Logging(日志):** 结构化日志,包含请求ID、节点名、输入输出摘要
代码示例:集成OpenTelemetry
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.instrumentation.requests import RequestsInstrumentor
# 初始化Tracer
provider = TracerProvider()
processor = BatchSpanProcessor(OTLPSpanExporter(endpoint="http://otel-collector:4318/v1/traces"))
provider.add_span_processor(processor)
trace.set_tracer_provider(provider)
tracer = trace.get_tracer(__name__)
# 自动追踪HTTP请求
RequestsInstrumentor().instrument()
# 在Agent节点中手动添加Span
def monitored_agent_node(state: dict) -> dict:
with tracer.start_as_current_span("intent_classification") as span:
span.set_attribute("user_query", state["user_query"][:100])
span.set_attribute("node_name", "classifier")
# 实际业务逻辑...
result = {"intent": "complaint"}
span.set_attribute("result", result["intent"])
return result
关键配置:Prometheus指标暴露
from prometheus_client import Counter, Histogram, generate_latest, start_http_server
import time
# 定义指标
pipeline_requests = Counter('pipeline_requests_total', 'Total pipeline requests', ['status'])
pipeline_duration = Histogram('pipeline_duration_seconds', 'Pipeline duration in seconds', buckets=[0.1, 0.5, 1, 2, 5, 10])
node_duration = Histogram('node_duration_seconds', 'Per-node duration', ['node_name'])
# 在流水线执行时记录
def run_with_metrics(pipeline_func):
@pipeline_duration.time()
def wrapper(state):
try:
result = pipeline_func(state)
pipeline_requests.labels(status='success').inc()
return result
except Exception as e:
pipeline_requests.labels(status='error').inc()
raise
return wrapper
# 启动HTTP服务暴露指标(默认端口8000)
start_http_server(8000)
---
三、实操步骤
步骤1:搭建最小流水线骨架
目标: 创建一个包含2个Agent节点的简单流水线,并成功运行
# 1. 创建项目目录
mkdir agent-pipeline && cd agent-pipeline
python -m venv venv
source venv/bin/activate # Windows: venv\Scripts\activate
# 2. 安装依赖
pip install langgraph redis prometheus-client opentelemetry-api opentelemetry-sdk
# 3. 创建主文件 main.py
main.py 核心代码:
from langgraph.graph import StateGraph, END
from typing import TypedDict
class State(TypedDict):
input_text: str
processed: str
def node_a(state: State):
print(f"Node A processing: {state['input_text']}")
return {"processed": f"A->{state['input_text']}"}
def node_b(state: State):
print(f"Node B processing: {state['processed']}")
return {"processed": f"B->{state['processed']}"}
builder = StateGraph(State)
builder.add_node("node_a", node_a)
builder.add_node("node_b", node_b)
builder.set_entry_point("node_a")
builder.add_edge("node_a", "node_b")
builder.add_edge("node_b", END)
app = builder.compile()
result = app.invoke({"input_text": "Hello Pipeline"})
print(f"Final result: {result}")
预期效果:
Node A processing: Hello Pipeline
Node B processing: A->Hello Pipeline
Final result: {'input_text': 'Hello Pipeline', 'processed': 'B->A->Hello Pipeline'}
步骤2:集成外部工具(数据库查询)
目标: 添加一个SQLite数据库查询工具,让Agent可以查询订单数据
# 1. 创建测试数据库
python -c "
import sqlite3
conn = sqlite3.connect('orders.db')
conn.execute('CREATE TABLE orders (id INT, product TEXT, amount REAL)')
conn.execute(\"INSERT INTO orders VALUES (1, 'Laptop', 999.99)\")
conn.execute(\"INSERT INTO orders VALUES (2, 'Mouse', 25.50)\")
conn.commit()
conn.close()
"
修改main.py,添加工具调用:
import sqlite3
def query_orders(state: State):
conn = sqlite3.connect('orders.db')
cursor = conn.cursor()
cursor.execute("SELECT * FROM orders")
rows = cursor.fetchall()
conn.close()
state["orders"] = rows
print(f"Found {len(rows)} orders")
return state
# 在构建图中添加节点
builder.add_node("query_orders", query_orders)
builder.add_edge("node_a", "query_orders")
builder.add_edge("query_orders", "node_b")
预期效果: 流水线会先处理文本,然后查询数据库,最后输出订单数据。
步骤3:配置人工审批节点
目标: 在关键节点插入暂停,等待Webhook恢复
# 1. 启动Redis
docker run -d -p 6379:6379 redis
# 2. 修改main.py,添加审批逻辑
审批节点代码:
import redis
import json
from datetime import datetime
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
def approval_node(state: State):
pipeline_id = f"pipeline_{datetime.now().timestamp()}"
state["_pipeline_id"] = pipeline_id
state["_status"] = "pending"
# 保存状态到Redis
r.set(f"pipeline:{pipeline_id}:state", json.dumps(state))
r.set(f"pipeline:{pipeline_id}:status", "pending_approval")
print(f"Pipeline {pipeline_id} 暂停,等待审批...")
print(f"请通过以下API恢复:POST /resume?pipeline_id={pipeline_id}&decision=approve")
# 模拟等待(实际应使用异步机制)
while True:
status = r.get(f"pipeline:{pipeline_id}:status")
if status == "resumed":
state_str = r.get(f"pipeline:{pipeline_id}:state")
return json.loads(state_str)
time.sleep(1)
恢复API(使用Flask):
from flask import Flask, request
app = Flask(__name__)
@app.route('/resume', methods=['POST'])
def resume():
pipeline_id = request.args.get('pipeline_id')
decision = request.args.get('decision', 'approve')
state_str = r.get(f"pipeline:{pipeline_id}:state")
if not state_str:
return {"error": "Pipeline not found"}, 404
state = json.loads(state_str)
state["_decision"] = decision
r.set(f"pipeline:{pipeline_id}:state", json.dumps(state))
r.set(f"pipeline:{pipeline_id}:status", "resumed")
return {"status": "resumed"}
if __name__ == '__main__':
app.run(port=5000)
预期效果: 流水线执行到审批节点会暂停,调用curl -X POST http://localhost:5000/resume?pipeline_id=xxx&decision=approve后继续执行。
---
四、常见问题与故障排查
问题