企业级Agent自动化流水线实战

🏷️ L4 📊 advanced ⏱️ 55分钟 🏷️ Agent流水线,企业级自动化,人工审批,可观测性,安全合规,前沿

# 企业级Agent自动化流水线实战

一、概述

为什么要学习这个主题

随着大模型技术的成熟,企业开始将AI Agent从“实验性玩具”转向“生产力工具”。但单点Agent在复杂业务场景中往往表现不佳:上下文丢失、工具调用混乱、缺乏安全管控。企业级Agent自动化流水线正是为了解决这些问题而生——它像一条工业生产线,将多个Agent、工具、审批节点、监控机制串联起来,形成稳定、可审计、可扩展的自动化系统。

行业背景: 2024-2025年,金融、医疗、电商等行业已出现大量Agent流水线案例,例如:

学完本课程能做什么

1. 设计多Agent协作架构,能根据业务需求拆分任务并编排流水线 2. 整合5种以上外部工具(数据库、API、文件系统、Webhook等) 3. 配置人工审批与异常回退机制,确保关键节点可控 4. 搭建监控与告警体系,实时掌握流水线运行状态 5. 优化流水线性能,处理高并发场景(如每秒100+请求)

适合人群与前置知识

---

二、核心知识点

模块1:流水线架构设计——从单体到编排

原理讲解

流水线本质是一个有向无环图(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后继续执行。

---

四、常见问题与故障排查

问题

在博海学习网开始学习 →