AI Agent 30天速成|Day10 笔记
今日核心主题:Agent 工具编排 + 异步批量任务 + 接口鉴权与分布式锁
修订说明
基于Day9全工具化架构迭代,保留Chroma向量库、统一工具网关、Redis记忆、熔断限流体系
新增工具链批量编排:支持顺序/并行多工具组合任务,无需模型多轮ReAct循环
实现异步后台任务队列,长耗时知识库批量入库、大规模向量离线计算不阻塞HTTP接口
增加全局接口Token鉴权、分布式会话锁,适配多实例分布式部署场景
新增Agent调用成本统计模块,统计LLM/Embedding调用次数、Token消耗、计费估算
今日总学习目标(全天8h分配)
理论学习:2.5h
分层代码编写调试:4h
面试复盘背诵:1.5h
掌握工具编排语法,定义串行/并行工具流水线,实现一次性多工具批量执行
基于Redis List实现轻量异步任务队列,处理长耗时离线向量任务
实现全局API鉴权、用户资源配额(LLM/Embedding每日调用上限)
分布式锁解决多实例Redis会话并发读写错乱问题
新增调用成本统计模块,全链路采集Token消耗与接口调用量
整合流水线编排、异步任务、鉴权配额、分布式锁,升级生产级Agent服务
一、核心理论教学笔记
1 工具流水线编排(Tool Pipeline)
1.1 单轮ReAct存在短板
ReAct每轮只能执行单个工具,多步骤任务需要多次LLM思考,消耗大量Token与接口次数;固定流程(先检索再计算、批量入库)重复循环效率极低。
1.2 两种编排模式
串行流水线(依赖型):上一个工具输出作为下一个工具输入,顺序执行
示例:rag_search → calculator(先查公式再代入计算)
并行流水线(无依赖):多个工具同时并发执行,结果统一汇总
示例:同时执行多条rag_search多维度知识库检索
1.3 流水线结构化定义
通过Pydantic定义流水线JSON结构,模型可一次性生成多工具任务,网关自动解析批量调度,仅消耗单次LLM思考。
流水线内置变量传递:前序工具输出可作为后续工具入参,无需重复拼接Prompt。
2 Redis异步任务队列(轻量版,无需RabbitMQ)
2.1 适用场景
批量文档入库、海量文本Embedding、大规模知识库更新等耗时操作,同步HTTP接口会超时。
2.2 队列执行流程
前端提交批量任务 → 接口校验权限配额 → 写入Redis任务队列,返回task_id
后台独立协程循环消费队列,执行向量入库/Embedding批量工具
任务进度、执行结果持久存入Redis Hash,用户通过task_id轮询查询状态
失败任务自动重试2次,最终失败写入死信队列人工排查
3 分布式配套能力
3.1 全局API鉴权
所有接口请求携带access_token,服务内存维护token-角色-配额映射;区分普通用户/管理员每日调用上限,超量直接拦截。
3.2 分布式会话锁
多服务实例同时收到同一用户多轮提问,并发读写Redis会话会造成消息错乱;基于Redis SETNX实现会话互斥锁,单会话同一时间仅允许一条对话执行。
3.3 调用成本统计
每条LLM、Embedding调用采集输入/输出Token,存入Redis计数器;按厂商单价估算消耗费用,支持按用户、会话、日期维度统计用量。
4 Day9→Day10 完整升级链路
用户HTTP请求携带token鉴权 → 校验用户每日调用配额
→ 获取分布式会话锁,防止并发乱序
→ 读取Redis分层记忆,自动摘要/滑动窗口压缩
→ 模型二选一:
简单任务:单工具ReAct循环执行
固定多步骤任务:生成流水线编排,网关批量并行/串行执行工具
→ 短任务同步返回;长批量任务推送Redis异步队列,返回task_id
→ 所有工具执行埋点:限流、熔断、日志、Token消耗统计
→ 工具结果回填上下文,LLM生成回答,脱敏处理
→ 对话持久Redis,释放分布式会话锁
→ 后台协程持续消费异步队列,处理离线向量批量任务
5 生产级配额与限流规则
访客:仅基础闲聊,禁用所有工具,每日LLM 10次上限
普通用户:calculator/rag_search可用,禁用vector_add;LLM 200次/天,Embedding 500次/天
管理员:全工具开放,无每日硬上限,仅全局令牌桶流量控制
二、今日学习重点
实现ToolPipeline流水线结构化定义,网关自动解析串行/并行批量工具执行
搭建Redis轻量异步任务队列,实现长耗时向量任务后台异步处理
开发全局Token鉴权、用户每日调用配额拦截逻辑
Redis分布式锁解决多实例会话并发读写冲突
全链路Token采集、用量统计、成本估算模块开发
兼容原有Day9所有能力,无破坏性改动,平滑升级
三、今日难点 & 解决方案
难点1 流水线工具之间参数传递复杂,前序结果无法自动传给下一个工具
解决方案:流水线结构定义输出变量key,网关缓存每一步工具结果,后续工具可通过${变量名}动态填充入参,自动替换文本。
难点2 多实例部署,同一用户并发对话导致历史消息错乱
解决方案:执行对话前后加分布式Redis锁,锁过期时间30s;对话流程未完成时,同会话新请求直接返回“对话处理中,请稍后”。
难点3 批量向量入库耗时过长,同步接口60s超时
解决方案:耗时任务剥离至Redis异步队列,接口仅生成任务ID立即返回;后台协程单独消费,前端轮询任务进度。
难点4 用户无感知超额调用LLM/Embedding产生高额费用
解决方案:按用户角色设置每日调用配额,Redis计数器实时累加,达到上限直接拦截,同时每日零点自动重置计数。
难点5 异步任务大量失败堆积无法感知
解决方案:区分正常队列与死信队列,失败重试2次仍异常的任务转入死信;日志记录全部失败task_id,提供查询接口查看失败详情。
四、完整项目代码(基于Day9扩展,新增文件+修改原有文件)
依赖新增
pip install python-multipart
项目目录新增文件
day10_agent/
├── .env # 新增鉴权、队列、配额配置
├── middleware.py # 新增分布式锁、用量统计逻辑
├── security.py # 新增token鉴权、用户配额校验
├── pipeline.py # 【今日新增】工具流水线编排定义
├── async_task.py # 【今日新增】Redis异步任务队列消费
├── llm_client.py
├── tool_gateway.py # 适配流水线批量工具执行
├── memory_store.py
├── agent_core.py # 支持流水线+ReAct双模式
└── main.py # 新增鉴权中间件、任务查询接口
1 .env 新增配置项
原有LLM/Redis/Chroma配置保留
API鉴权
ADMIN_TOKEN=admin123456
USER_TOKEN=user654321
GUEST_TOKEN=guest0000
每日调用配额
ADMIN_QUOTA_LLM=9999
ADMIN_QUOTA_EMB=9999
USER_QUOTA_LLM=200
USER_QUOTA_EMB=500
GUEST_QUOTA_LLM=10
GUEST_QUOTA_EMB=0
异步任务队列
TASK_QUEUE_KEY=agent:task:queue
DEAD_QUEUE_KEY=agent:task:dead
TASK_EXPIRE=86400
分布式锁
LOCK_PREFIX=agent:session:lock
LOCK_EXPIRE=30
Token计费单价(元/千token)
PRICE_INPUT=0.5
PRICE_OUTPUT=1.0
PRICE_EMB=0.2
2 pipeline.py(全新文件:工具流水线编排)
from pydantic import BaseModel, Field
from typing import List, Dict, Optional, Union
from tool_gateway import gateway
class PipelineStep(BaseModel):
task_id: str = Field(description=“流水线内步骤唯一标识”)
tool_name: str = Field(description=“执行工具名称”)
args: Dict = Field(description=“工具入参,支持${变量名}引用上一步输出”)
depends_on: List[str] = Field(default=[], description=“依赖步骤id,空则并行执行”)
output_var: str = Field(description=“本步骤输出存储变量名,供后续步骤引用”)
class ToolPipeline(BaseModel):
run_mode: str = Field(description=“serial串行 / parallel并行”)
steps: List[PipelineStep] = Field(description=“流水线步骤列表”)
class PipelineRunner:
def init(self):
self.var_cache = {} # 存储各步骤输出变量
# 替换参数内${xxx}变量
def replace_var(self, raw_arg: Union[str, Dict]) -> Union[str, Dict]:
if isinstance(raw_arg, str):
for k, v in self.var_cache.items():
raw_arg = raw_arg.replace(f"${{{k}}}", str(v))
return raw_arg
if isinstance(raw_arg, dict):
new_dict = {}
for k, val in raw_arg.items():
new_dict[k] = self.replace_var(val)
return new_dict
return raw_arg
async def run_single_step(self, step: PipelineStep, trace_id: str, user_role: str):
parsed_args = self.replace_var(step.args)
res = await gateway.run_tool(step.tool_name, parsed_args, trace_id, user_role)
self.var_cache[step.output_var] = res
return {step.task_id: res}
async def run_pipeline(self, pipeline: ToolPipeline, trace_id: str, user_role: str):
all_results = {}
step_map = {s.task_id: s for s in pipeline.steps}
finished = set()
while len(finished) < len(step_map):
run_coros = []
run_steps = []
for sid, step in step_map.items():
if sid in finished:
continue
# 判断依赖是否全部完成
dep_ok = all(d in finished for d in step.depends_on)
if dep_ok:
run_coros.append(self.run_single_step(step, trace_id, user_role))
run_steps.append(sid)
if not run_coros:
break
batch_res = await asyncio.gather(*run_coros)
for item in batch_res:
all_results.update(item)
for s in run_steps:
finished.add(s)
return {
"pipeline_var": self.var_cache,
"step_result": all_results
}
pipeline_runner = PipelineRunner()
更多推荐



所有评论(0)