)
import asyncioimport jsonimport uuidimport osfrom dotenv import load_dotenvfrom memory_store import memoryfrom pipeline import ToolPipeline, pipeline_runnerload_dotenv()TASK_QUEUE os.getenv(“TASK_QUEUE_KEY”)DEAD_QUEUE os.getenv(“DEAD_QUEUE_KEY”)TASK_EXPIRE int(os.getenv(“TASK_EXPIRE”))class AsyncTaskCenter:definit(self):self.redis memory.redisdef gen_task_id(self): return str(uuid.uuid4()) async def submit_task(self, task_type: str, pipeline_dict: dict, trace_id: str, user_role: str): tid self.gen_task_id() task_data json.dumps({ task_id: tid, task_type: task_type, pipeline: pipeline_dict, trace_id: trace_id, user_role: user_role, status: pending, retry_cnt: 0 }) # 写入队列 Hash存储详情 await self.redis.lpush(TASK_QUEUE, task_data) await self.redis.setex(ftask:info:{tid}, TASK_EXPIRE, task_data) return tid async def get_task_info(self, task_id: str): raw await self.redis.get(ftask:info:{task_id}) if not raw: return None return json.loads(raw) async def consumer_loop(self): print(异步任务消费协程启动成功) while True: # 阻塞读取队列 raw_task await self.redis.brpop(TASK_QUEUE, timeout2) if not raw_task: continue _, task_str raw_task task json.loads(task_str) tid task[task_id] try: task[status] running await self.redis.setex(ftask:info:{tid}, TASK_EXPIRE, json.dumps(task)) # 执行流水线 pipeline ToolPipeline(**task[pipeline]) res await pipeline_runner.run_pipeline(pipeline, task[trace_id], task[user_role]) task[status] success task[result] json.dumps(res, ensure_asciiFalse) except Exception as e: task[retry_cnt] 1 task[err_msg] str(e) if task[retry_cnt] 2: # 重新入队重试 await self.redis.lpush(TASK_QUEUE, json.dumps(task)) continue task[status] failed # 转入死信队列 await self.redis.rpush(DEAD_QUEUE, json.dumps(task)) await self.redis.setex(ftask:info:{tid}, TASK_EXPIRE, json.dumps(task))task_center AsyncTaskCenter()4 middleware.py 新增分布式锁、用量统计代码在原有代码末尾追加import osLOCK_PREFIX os.getenv(“LOCK_PREFIX”)LOCK_EXPIRE int(os.getenv(“LOCK_EXPIRE”))PRICE_INPUT float(os.getenv(“PRICE_INPUT”))PRICE_OUTPUT float(os.getenv(“PRICE_OUTPUT”))PRICE_EMB float(os.getenv(“PRICE_EMB”))分布式会话锁async def session_lock(redis, session_id: str):key f{LOCK_PREFIX}{session_id}ok await redis.set(key, “locked”, exLOCK_EXPIRE, nxTrue)return bool(ok)async def session_unlock(redis, session_id: str):key f{LOCK_PREFIX}{session_id}await redis.delete(key)Token用量统计class TokenStat:definit(self, redis):self.redis redisasync def add_llm_token(self, user_role: str, input_tok: int, output_tok: int): today time.strftime(%Y%m%d) key_llm_in fstat:{today}:{user_role}:llm_input key_llm_out fstat:{today}:{user_role}:llm_output await self.redis.incrby(key_llm_in, input_tok) await self.redis.incrby(key_llm_out, output_tok) async def add_emb_token(self, user_role: str, tok_num: int): today time.strftime(%Y%m%d) key_emb fstat:{today}:{user_role}:embedding await self.redis.incrby(key_emb, tok_num) async def get_cost(self, user_role: str): today time.strftime(%Y%m%d) in_tok int(await self.redis.get(fstat:{today}:{user_role}:llm_in) or 0) out_tok int(await self.redis.get(fstat:{today}:llm_out) or 0) emb_tok int(await self.redis.get(fstat:{today}:emb) or 0) cost (in_tok / 1000) * PRICE_INPUT (out_tok / 1000) * PRICE_OUTPUT (emb_tok / 1000) * PRICE_EMB return round(cost, 4)stat_client TokenStat(memory.redis)5 security.py 新增Token鉴权配额校验原有代码保留追加下方内容import osTOKEN_ROLE_MAP {os.getenv(“ADMIN_TOKEN”): “admin”,os.getenv(“USER_TOKEN”): “user”,os.getenv(“GUEST_TOKEN”): “guest”}QUOTA_CFG {“admin”: {“llm”: int(os.getenv(“ADMIN_QUOTA_LLM”)),“emb”: int(os.getenv(“ADMIN_QUOTA_EMB”))},“user”: {“llm”: int(os.getenv(“USER_QUOTA_LLM”)),“emb”: int(os.getenv(“USER_QUOTA_EMB”))},“guest”: {“llm”: int(os.getenv(“GUEST_QUOTA_LLM”)),“emb”: int(os.getenv(“GUEST_QUOTA_EMB”))}}Token鉴权def verify_token(token: str) - tuple[bool, str]:if token not in TOKEN_ROLE_MAP:return False, “”return True, TOKEN_ROLE_MAP[token]校验当日调用配额async def check_quota(redis, user_role: str, call_type: str) - bool:today time.strftime(“%Y%m%d”)limit QUOTA_CFG[user_role][call_type]key fquota:{today}:{user_role}:{call_type}current int(await redis.get(key) or 0)if current limit:return Falseawait redis.incr(key)return True6 agent_core.py 改造支持流水线编排原有导入追加from pipeline import ToolPipeline, pipeline_runnerfrom middleware import stat_clientclass ReActAgent:# 原有init、reflect函数不变async def run_chat(self, session_id: str, user_role: str, user_input: str, trace_id: str):history await memory.load_history(session_id)base_msg [{“role”:“system”, “content”:self.system_prompt}] historybase_msg.append({“role”:“user”, “content”:user_input})msg_list await memory.auto_compress(base_msg)# 新增判断生成流水线还是单轮ReActpipeline_prompt “”判断用户问题是否为固定多步骤任务是则输出ToolPipeline JSONrun_mode选serial/parallel简单单任务输出{“is_pipeline”:false}“”pipe_check_msg msg_list [{“role”:“user”, “content”:pipeline_prompt}]pipe_raw await llm_client.chat_sync(pipe_check_msg, temperature0.0)[“choices”][0][“message”][“content”]# 省略JSON解析逻辑区分两种执行分支if “is_pipeline” in pipe_raw and json.loads(pipe_raw)[“is_pipeline”] is False:# 原有ReAct循环逻辑不变执行后统计tokenloop_cnt 0tool_record {}input_tokens sum(len(m[“content”]) for m in msg_list)while loop_cnt self.max_loop:# 原有ReAct代码不变…output_tokens len(final_raw)await stat_client.add_llm_token(user_role, input_tokens, output_tokens)else:# 流水线分支pipe_data json.loads(pipe_raw)pipeline ToolPipeline(**pipe_data)pipeline_res await pipeline_runner.run_pipeline(pipeline, trace_id, user_role)concat_info json.dumps(pipeline_res[“step_result”], ensure_asciiFalse)final_prompt msg_list [{“role”:“user”, “content”:f结合流水线结果回答{concat_info}}]final_raw await llm_client.chat_sync(final_prompt, temperature0.1)[“choices”][0][“message”][“content”]input_tokens sum(len(m[“content”]) for m in final_prompt)output_tokens len(final_raw)await stat_client.add_llm_token(user_role, input_tokens, output_tokens)# 脱敏、持久化逻辑不变final_ans desensitize(final_raw)await memory.append(session_id, “user”, user_input)await memory.append(session_id, “assistant”, final_ans)log_client.write(trace_id, “INFO”, {“final_answer”: final_ans[:300]})return {“trace_id”: trace_id,“tool_record”: tool_record if “tool_record” in locals() else pipeline_res,“answer”: final_ans}react_agent ReActAgent()7 main.py 新增鉴权、任务查询、异步提交接口from fastapi import FastAPI, Query, Header, HTTPExceptionimport asynciofrom agent_core import react_agentfrom memory_store import memoryfrom middleware import global_bucket, create_trace_id, log_client, session_lock, session_unlockfrom security import input_verify, verify_token, check_quotafrom async_task import task_centerapp FastAPI(title“Day10 流水线Agent异步任务分布式锁鉴权配额”)全局启动事件新增异步消费协程app.on_event(“startup”)async def startup():await memory.connect()asyncio.create_task(task_center.consumer_loop())app.on_event(“shutdown”)async def shutdown():await memory.close()通用对话接口增加token鉴权app.get(“/agent/chat”)async def chat_api(Authorization: str Header(…, description“Bearer token”),session_id: str Query(…),user_role: str Query(default“user”),prompt: str Query(…)):# 1 Token鉴权token Authorization.replace(“Bearer “,””)auth_ok, real_role verify_token(token)if not auth_ok:raise HTTPException(status_code401, detail“非法访问令牌”)# 2 配额校验llm_quota_ok await check_quota(memory.redis, real_role, “llm”)if not llm_quota_ok:return {“trace_id”: create_trace_id(), “answer”: “今日LLM调用额度已耗尽请明日再试”}# 3 分布式会话锁lock_success await session_lock(memory.redis, session_id)if not lock_success:return {“trace_id”: create_trace_id(), “answer”: “当前会话正在处理对话请稍后重试”}try:trace_id create_trace_id()log_client.write(trace_id, “INFO”, {“session_id”: session_id, “input”: prompt, “role”: real_role})# 输入安全校验ok, safe_text await input_verify(prompt)if not ok:return {“trace_id”: trace_id, “answer”: safe_text}# 全局限流if not await global_b.get_token():log_client.write(trace_id, “WARN”, {“msg”: “触发全局限流”})return {“trace_id”: trace_id, “answer”: “服务繁忙请稍后重试”}# 主Agent执行result await react_agent.run_chat(session_id, real_role, safe_text, trace_id)return resultfinally:# 释放会话锁await session_unlock(memory.redis, session_id)提交异步批量流水线任务接口app.post(“/agent/task/submit”)async def submit_task(Authorization: str Header(…),session_id: str,task_type: str,pipeline: dict):token Authorization.replace(“Bearer “,””)auth_ok, real_role verify_token(token)if not auth_ok:raise HTTPException(401, “token无效”)# 仅管理员可提交批量入库任务if task_type “batch_embed” and real_role ! “admin”:raise HTTPException(403, “仅管理员可执行批量向量任务”)trace_id create_trace_id()tid await task_center.submit_task(task_type, pipeline, trace_id, real_role)return {“task_id”: tid, “trace_id”: trace_id, “msg”: “任务已进入后台队列”}查询异步任务进度app.get(“/agent/task/query”)async def query_task(task_id: str):info await task_center.get_task_info(task_id)if not info:return {“msg”: “任务不存在或已过期”}return infoifname “main”:import uvicornuvicorn.run(“main.py”, reloadTrue)五、今日实操练习任务配置.env内三类Token分别用admin/user/guest令牌调用接口验证权限隔离与每日配额拦截构造复合问题先检索RAG定义再计算(10020)*5观察模型自动生成串行流水线批量执行工具使用admin账号提交批量文档入库流水线获取task_id轮询/agent/task/query查看后台执行进度同一session_id并发发送两条提问验证分布式锁拦截返回会话繁忙提示连续调用接口耗尽user角色每日LLM配额观察配额拦截提示人为制造批量向量任务异常查看死信队列存储失败任务信息查看日志与Redis统计key验证LLM/Embedding Token消耗与费用估算正常六、配套完整面试题含标准答案基础问答1 Tool流水线和ReAct循环的核心区别各自适用场景标准答案ReAct单次只能调用一个工具多步骤任务需要多轮LLM思考消耗更多Token适合无固定流程、开放式未知问题。Tool流水线一次性定义串行/并行多工具步骤模型仅一次规划网关批量调度适合流程固定、多依赖/多并行查询的标准化任务批量知识库入库、固定计算检索组合优先使用流水线。2 为什么要用Redis实现异步任务不直接同步执行批量向量入库标准答案批量Embedding、大规模文档入库耗时可达数十秒HTTP接口默认超时会断开连接异步队列将耗时操作剥离后台前端立刻获取task_id轮询结果提升接口可用性同时失败任务支持自动重试不阻塞用户对话主线程。3 分布式会话锁解决什么问题实现原理标准答案多服务实例部署时同一用户同时发起多条对话请求并发读写Redis会话列表会造成消息错乱、上下文拼接异常。基于Redis SETNX原子命令加锁会话处理期间持有30s过期锁新同会话请求检测锁存在直接拒绝对话结束主动释放锁保证单会话串行执行。4 用户配额统计基于Redis什么数据结构零点如何重置每日用量标准答案使用Redis String计数器key拼接日期用户角色区分每日用量每日零点定时任务删除当日统计key自动重置次日计数设置过期时间24h自动清理过期历史用量数据。5 流水线变量传递实现逻辑是什么标准答案网关内置变量缓存字典每一步工具执行完成后将输出存入对应output_var解析后续步骤入参时正则匹配${变量名}自动替换为上一步工具输出无需模型重复拼接上下文参数。工程实操题1 流水线步骤存在循环依赖A依赖BB依赖A如何处理标准答案规划阶段Prompt增加依赖校验规则模型生成流水线时禁止循环依赖网关执行前遍历所有depends_on做拓扑校验检测循环依赖直接返回任务规划失败拒绝执行。2 异步任务消费协程意外退出如何保证任务不丢失标准答案Redis brpop阻塞读取特性任务取出未处理完成进程崩溃会重新回到队列可搭配定时补偿协程扫描长时间running状态未完成的任务重新入队保障消息可靠消费。3 多角色权限双层配额工具权限每日调用上限如何分层校验标准答案第一层security工具权限guest禁用所有工具user禁用vector_add第二层每日用量配额Redis计数器拦截LLM/Embedding超量请求两层校验全部通过才允许执行对应工具任一拦截直接返回提示。4 线上如何防止恶意刷取Embedding产生高额费用标准答案三层防护①全局令牌桶流量限流②按用户角色设置每日Embedding调用硬配额③异步批量向量任务仅管理员开放普通用户无法批量发起大规模Embedding请求。