ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

AI应用后端实战:持久化、HITL、流式输出与安全架构设计

AI应用后端实战:持久化、HITL、流式输出与安全架构设计 1. 项目概述从零构建一个健壮的AI应用后端最近在做一个AI驱动的对话应用后端从零开始踩了不少坑也积累了一些心得。这个项目不是简单的调用API而是涉及到了智能体Agent的复杂编排、长对话管理、以及如何安全稳定地对外提供服务。整个过程大概持续了8天标题里的第16~23天核心聚焦在几个关键问题上如何让智能体的“记忆”不丢失持久化如何让人类在关键时刻介入决策HITL如何让用户获得流畅的交互体验流式输出如何安全地扩展智能体的能力MCP以及如何确保整个系统不被滥用安全防护这五个点任何一个没处理好都会让应用变得难用、不稳定甚至危险。比如用户聊了十分钟刷新页面后AI就失忆了或者一个自动化的智能体做出了错误的财务操作又或者接口被恶意刷量导致服务瘫痪。我在这段时间里主要就是和这些问题“搏斗”最终形成了一个相对完整的后端解决方案。如果你也在开发类似的AI应用尤其是涉及复杂Agent和实时交互的场景希望我的这些实践和思考能给你一些直接的参考。2. 持久化让智能体的“记忆”落地生根智能体Agent的核心价值在于其上下文感知和持续学习的能力。但这一切的前提是它的“状态”和“记忆”必须能够被可靠地保存下来。否则每次服务重启、会话中断智能体都会“失忆”用户体验会大打折扣。2.1 为什么需要持久化持久化不仅仅是把聊天记录存到数据库那么简单。对于一个智能体系统我们需要持久化的对象至少包括会话历史用户与智能体之间的完整对话记录。智能体状态Agent在运行过程中的内部状态例如它当前的目标、已执行的动作列表、工具调用的结果缓存等。用户配置与偏好用户对智能体的个性化设置。长期记忆从历史对话中提炼出的、关于用户或领域的关键知识。如果这些数据只存在于内存中那么服务重启即丢失任何运维操作、版本更新、甚至意外的进程崩溃都会导致数据清零。无法水平扩展在多实例部署时用户的请求可能被路由到不同的服务器实例如果状态没有集中存储会话将无法保持连贯。无法进行离线分析与再训练历史数据是优化智能体表现、进行数据分析的宝贵资产。2.2 技术选型与实战Redis 关系型数据库我采用的是一种混合持久化策略核心是Redis和PostgreSQL或其他关系型数据库的组合。Redis负责高速缓存和会话状态存储。它的内存读写性能极高非常适合存储需要快速访问的临时状态和热数据。实战配置以Docker为例 很多人用Docker跑Redis时会发现重启容器后数据没了这是因为默认配置没有开启持久化。必须修改redis.conf或通过命令参数启用。# docker-compose.yml 示例片段 services: redis: image: redis:7-alpine command: redis-server --appendonly yes --appendfsync everysec volumes: - ./redis-data:/data - ./config/redis.conf:/usr/local/etc/redis/redis.conf # 挂载自定义配置文件 ports: - 6379:6379--appendonly yes启用AOFAppend Only File持久化记录每一个写操作。--appendfsync everysec每秒同步一次AOF文件到磁盘在性能和数据安全间取得平衡。你也可以设为always最安全但慢或no由操作系统决定不安全。挂载/data卷是为了让AOF文件或RDB快照文件在容器外持久化。数据结构设计会话历史使用List类型存储Key为chat:session:{session_id}:messagesValue为序列化的消息对象列表。List的LPUSH/RPUSH和LRANGE操作非常适合消息的追加和范围读取。智能体状态使用Hash类型存储Key为agent:state:{session_id}Field-Value对存储各种状态属性。Hash适合存储一个对象的多个字段。限流与计数器使用String类型配合INCR和EXPIRE命令实现简单的API调用频率限制。PostgreSQL负责永久存储和复杂查询。关系型数据库适合存储需要长期保留、结构复杂、且可能需要复杂关联查询的数据。表结构设计举例-- 用户表 CREATE TABLE users ( id UUID PRIMARY KEY, email VARCHAR(255) UNIQUE, created_at TIMESTAMPTZ DEFAULT NOW() ); -- 会话表与用户关联 CREATE TABLE sessions ( id UUID PRIMARY KEY, user_id UUID REFERENCES users(id) ON DELETE CASCADE, title VARCHAR(255), -- 会话标题可由AI自动生成 created_at TIMESTAMPTZ DEFAULT NOW(), updated_at TIMESTAMPTZ DEFAULT NOW() ); -- 消息表与会话关联 CREATE TABLE messages ( id BIGSERIAL PRIMARY KEY, session_id UUID REFERENCES sessions(id) ON DELETE CASCADE, role VARCHAR(20) NOT NULL, -- user, assistant, system, tool content TEXT NOT NULL, tool_calls JSONB, -- 存储工具调用请求如果是assistant消息 tool_call_id VARCHAR(255), -- 工具调用ID如果是tool消息 created_at TIMESTAMPTZ DEFAULT NOW() ); -- 长期记忆/知识片段表 CREATE TABLE memories ( id BIGSERIAL PRIMARY KEY, user_id UUID REFERENCES users(id) ON DELETE CASCADE, content TEXT NOT NULL, -- 记忆内容 embedding vector(1536), -- 向量化表示如果使用向量搜索 metadata JSONB, -- 来源、类型等元数据 created_at TIMESTAMPTZ DEFAULT NOW() );注意JSONB类型是PostgreSQL的强大功能可以灵活存储半结构化的数据如工具调用参数、元数据等同时支持索引和查询。2.3 数据同步与一致性策略混合存储带来了数据一致性的挑战。我的策略是写数据库刷缓存。消息写入流程用户发送消息后首先将其写入PostgreSQL的messages表。写入成功后立即将该消息RPUSH到Redis中对应的会话消息列表。更新Redis中会话的updated_at时间戳或一个简单的过期时间续期。消息读取流程首先尝试从Redis中获取指定范围的会话历史LRANGE。如果Redis中没有例如首次查询或缓存过期则从PostgreSQL中查询并回写到Redis中。状态同步智能体的运行状态如计划、中间结果主要存在Redis中保证高速读写。当会话结束时或定期如每10分钟将重要的最终状态快照持久化到PostgreSQL的一个agent_snapshots表中用于审计和恢复。这种策略保证了最新数据的快速访问同时确保了数据的最终可靠性。对于关键操作如创建订单、修改数据库需要更严格的事务性则应直接在数据库层面完成。3. HITL为智能体装上“紧急制动”和“方向盘”HITLHuman-in-the-Loop人在回路是确保AI应用安全、可靠、可控的关键机制。我们不能完全放任智能体自动化执行所有操作尤其是在涉及敏感操作、重大决策或模糊边界时。3.1 HITL的典型场景关键操作审批智能体试图执行“发送邮件”、“创建订单”、“支付”、“修改生产数据库”等操作前需要暂停并等待用户或管理员在界面上点击“确认”。模糊意图澄清当用户指令不明确智能体无法通过简单追问确定时可以生成几个可能的选项让用户选择。结果审核与修正智能体生成了一段文案、一份报告或一段代码提交给用户做最终审核和编辑。异常处理当工具调用失败、返回意外结果或触发安全规则时转入人工处理流程。3.2 实现方案状态机与任务队列实现HITL的核心是设计一个状态机来管理智能体的执行流程。定义智能体状态RUNNING: 正常执行中。AWAITING_APPROVAL: 已暂停等待用户审批某个待执行的动作Tool Call。AWAITING_CLARIFICATION: 已暂停等待用户澄清意图或做出选择。PAUSED_BY_USER: 被用户手动暂停。COMPLETED: 任务完成。FAILED: 任务失败。修改智能体执行循环 在智能体的主循环中在执行任何一个“敏感工具”之前插入检查逻辑。# 伪代码示例 class AgentWithHITL: async def run(self, user_input): self.state RUNNING # ... 规划、思考 ... for action in planned_actions: if self._requires_approval(action): # 1. 暂停智能体将状态置为 AWAITING_APPROVAL self.state AWAITING_APPROVAL self.pending_action action # 2. 将待审批事件持久化到数据库 approval_task_id await self._save_approval_task(action) # 3. 通过WebSocket或轮询API通知前端 await self._notify_frontend_approval_required(approval_task_id, action) # 4. 等待异步直到收到前端响应或超时 user_decision await self._wait_for_user_decision(approval_task_id, timeout300) if user_decision APPROVED: self.state RUNNING result await self._execute_action(action) elif user_decision REJECTED: # 用户拒绝智能体需要调整计划 self.state RUNNING self._replan_without_action(action) continue else: # TIMEOUT self.state FAILED raise TimeoutError(用户审批超时) else: # 无需审批直接执行 result await self._execute_action(action) # ... 处理结果继续循环 ...前后端协作后端提供API端点用于前端列出待审批任务、提交审批决定通过/拒绝/修改参数。前端在聊天界面中当收到“等待审批”的事件时渲染一个突出的UI组件例如一个模态框或一个固定在输入框上方的审批栏展示待执行的动作详情并提供“批准”和“拒绝”按钮。通信使用WebSocket可以实现最实时的双向通信。后端在需要审批时推送事件前端在用户操作后推送决策。退而求其次可以使用前端短轮询Polling一个“待处理任务”API。任务队列的引入 对于需要长时间等待人工审批的场景智能体的执行进程不能一直阻塞。此时可以引入任务队列如Celery Redis/RabbitMQ。当遇到需要HITL的节点时智能体任务被挂起一个“人工审批任务”被创建并存入数据库。一个独立的Worker进程监听“人工任务完成”的事件。用户在界面上完成审批后系统发布事件Worker接收到事件后根据任务ID找到对应的智能体执行上下文并唤醒其继续执行。实操心得HITL的粒度设计很重要。不要事无巨细都让人审批那会严重损害体验。应该基于“风险”和“成本”来定义审批规则。例如通过一个配置文件来声明哪些工具Tool需要审批甚至可以配置不同参数组合下的不同审批级别。同时一定要设置超时时间防止任务因无人处理而永远挂起。4. 流式输出打造“正在思考”的实时体验流式输出Streaming对于AI对话应用几乎是必需品。它能极大提升用户体验让用户感觉AI在实时“思考”和“回应”而不是在长时间等待后突然蹦出一大段文字。4.1 技术选型SSE vs. WebSocket实现流式主要有两种协议Server-Sent Events (SSE)和WebSocket。SSE优点基于HTTP协议简单天然支持自动重连浏览器端有EventSourceAPI直接支持。它是单向的服务器到客户端非常适合只需服务器推送的场景如聊天回复、通知。缺点单向通信如果对话中需要双向实时交互如一边流式输出一边随时打断则力有不逮。WebSocket优点全双工通信功能强大可以处理任何需要双向实时数据的场景。缺点协议相对复杂需要自己处理连接管理、心跳、重连等。我的选择是对于纯AI文本回复流式输出优先使用SSE对于需要复杂双向交互的智能体应用如需要随时发送工具调用结果、打断生成使用WebSocket。4.2 基于SSE的流式输出实现Python FastAPI示例以下是一个使用FastAPI和SSE实现流式回复的清晰示例# app.py from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import asyncio import json from typing import AsyncGenerator app FastAPI() async def simulate_llm_streaming(prompt: str) - AsyncGenerator[str, None]: 模拟大语言模型流式生成。 在实际项目中这里会调用OpenAI、Anthropic或本地模型的流式API。 # 假设模型返回的是一系列“词块” simulated_chunks [ 你好, 我是AI助手。, 你刚才问的问题是, prompt, 。, 这是一个很好的问题, 让我来详细解答一下... ] for chunk in simulated_chunks: # 模拟一点延迟更像真实生成过程 await asyncio.sleep(0.1) # 关键以SSE格式 yield 数据 # 格式为 data: 内容\n\n yield fdata: {json.dumps({content: chunk})}\n\n # 可以发送一个特殊事件表示结束例如 [DONE] yield fdata: {json.dumps({content: [DONE]})}\n\n app.post(/chat/stream) async def chat_stream(request: Request): data await request.json() prompt data.get(prompt, ) # 使用StreamingResponse并设置正确的媒体类型 return StreamingResponse( simulate_llm_streaming(prompt), media_typetext/event-stream, # SSE的MIME类型 headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, # 针对Nginx代理的重要设置禁用缓冲 } )前端Vue/JavaScript如何接收// 前端代码示例 async function streamChat(prompt) { const eventSource new EventSource(/chat/stream?prompt${encodeURIComponent(prompt)}); // 注意GET请求通常有URL长度限制复杂数据建议用POST EventSource polyfill 或直接使用fetch // 更推荐使用 fetch 来 POST 数据并读取流 const response await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt: prompt }) }); const reader response.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const lines buffer.split(\n\n); buffer lines.pop(); // 最后一条可能不完整放回buffer for (const line of lines) { if (line.startsWith(data: )) { const data line.slice(6); // 去掉 data: try { const parsed JSON.parse(data); if (parsed.content [DONE]) { console.log(流式传输结束); return; } // 将 parsed.content 追加到UI上 appendToChatUI(parsed.content); } catch (e) { console.error(解析SSE数据失败:, e, data); } } } } }重要提示在使用Nginx等反向代理时默认会对响应进行缓冲这会导致SSE流式输出有延迟甚至无法工作。必须在Nginx配置中为对应的location添加proxy_buffering off;指令。4.3 结合智能体的复杂流式当智能体在运行中调用工具时流式输出会变得更复杂。你不仅需要流式输出思考文本还需要输出工具调用的事件和结果。这时可以定义一套简单的事件协议{event: thinking, data: 我正在分析你的问题...} {event: tool_call, data: {name: search_web, args: {query: ...}}} {event: tool_result, data: {result: ...}} {event: message, data: {content: 根据搜索结果...}}前端根据event字段来渲染不同的UI元素如思考指示器、工具调用卡片、最终回答。5. MCP安全地扩展智能体的“技能包”MCPModel Context Protocol是一个新兴的协议旨在标准化AI模型尤其是智能体与外部工具、数据源之间的交互方式。你可以把它理解为智能体的“插件系统”或“驱动协议”。它解决了智能体如何安全、可控、动态地获取和使用外部能力的问题。5.1 MCP的核心思想在没有MCP之前我们通常需要硬编码工具列表到智能体系统中。每增加一个新工具如查询数据库、调用内部API、读取特定文件都需要修改后端代码重新部署。MCP试图将工具提供者Server和工具消费者Client即智能体运行时解耦。MCP Server提供一系列“工具”Tools和“资源”Resources。例如一个“数据库MCP Server”可以提供run_query工具一个“文件系统MCP Server”可以提供read_file、list_directory工具。MCP Client智能体运行时。它通过标准的MCP协议与一个或多个MCP Server通信动态地发现并调用这些工具。这样做的好处是动态能力发现智能体启动后可以自动发现当前可用的所有工具无需重启。安全边界清晰每个MCP Server可以运行在独立的、权限受控的沙箱环境中。即使某个Server被恶意利用或出现故障影响范围也被隔离。开发与集成标准化工具开发者只需遵循MCP协议实现Server就可以让所有兼容MCP的智能体如Claude Desktop、Cursor等使用。5.2 一个简单的MCP Server示例概念虽然MCP协议本身有详细规范但其核心是Server向Client宣告自己有哪些能力。下面是一个极度简化的概念性示例展示其思想假设我们用Python为一个“计算器”功能实现一个MCP Server。# 伪代码/概念展示 class CalculatorMCPServer: def __init__(self): self.tools [ { name: add, description: Add two numbers, inputSchema: { type: object, properties: { a: {type: number}, b: {type: number} }, required: [a, b] } }, { name: multiply, description: Multiply two numbers, inputSchema: {...} } ] async def handle_request(self, request): if request.method GET and request.path /tools: return json.dumps({tools: self.tools}) elif request.method POST and request.path /execute: data await request.json() tool_name data[name] arguments data[arguments] if tool_name add: result arguments[a] arguments[b] return json.dumps({content: [{type: text, text: str(result)}]}) # ... 处理其他工具智能体运行时MCP Client会连接到这个Server的地址例如http://localhost:8080通过/tools端点获取工具列表。当用户要求“计算一下123456”时智能体会决定调用add工具并向Server的/execute端点发送{name: add, arguments: {a: 123, b: 456}}最后将结果579整合到回复中。5.3 在项目中集成MCP目前直接从头实现MCP Server和Client需要处理协议细节如Stdio/SSE传输、初始化握手、资源管理等。更实用的方式是使用官方或社区的SDK。对于智能体框架Client侧如LangChain、LlamaIndex等已经开始集成MCP支持。你可以配置一个MCP Server的URL框架会自动发现并加载其工具。对于开发工具Server侧可以使用modelcontextprotocol/sdkNode.js或其他语言的SDK来快速构建MCP Server暴露你的内部API或数据库能力。集成MCP的关键考量身份验证与授权MCP Server暴露了能力必须确保只有受信的Client你的智能体后端可以连接。需要在Server端实现API密钥、JWT令牌或网络层如白名单的认证。输入验证与清理Server端必须对Client传入的参数进行严格的验证和清理防止注入攻击尤其是SQL、命令注入。资源访问控制不同的MCP Server应具有最小权限。例如一个“用户数据查询Server”不应该有删除权限。监控与日志所有工具调用必须有详细的日志记录便于审计和问题排查。6. 安全防护构建AI应用的后端防线AI应用的后端面临着独特的安全挑战高计算成本的接口容易被滥用导致巨额账单用户输入可能包含恶意指令Prompt Injection智能体可能被诱导执行危险操作。6.1 多层次防护策略我构建的安全防线主要包含以下几个层面1. 接入层防护网络与基础设施DDoS防护与速率限制使用云服务商如AWS WAF、Cloudflare或Nginx的limit_req模块在网关层对IP和用户ID实施严格的请求速率限制。对于/chat和/generate这类高成本端点限制必须更严格。Web应用防火墙WAF配置WAF规则过滤常见的SQL注入、XSS等攻击 payload即使它们藏在用户的聊天文本里。“本网站使用安全服务防护恶意自动程序”对于公开的、重要的端点可以引入人机验证如CAPTCHA或行为分析服务在检测到可疑自动化流量时例如来自数据中心的IP、极高的请求频率、无规律的鼠标移动弹出验证防止爬虫和恶意机器人滥用API。这行提示语常见于Cloudflare等安全服务。2. 应用层防护业务逻辑用户认证与授权所有API必须要求有效的身份令牌JWT。实施基于角色的访问控制RBAC确保用户只能访问自己被授权的资源和功能。输入净化与验证结构化数据对所有API参数使用强类型验证如Pydantic。非结构化文本用户输入建立“提示词防火墙”。这包括关键词过滤维护一个动态更新的黑名单过滤明显恶意、违法或涉及敏感话题的词汇。意图分类用一个轻量级模型对用户输入进行实时分类判断其是否为正常对话、尝试越权指令如“忽略之前的指令”、或攻击尝试。上下文长度限制防止用户通过输入超长文本来进行资源耗尽攻击。输出过滤与审查对AI生成的内容进行安全审查防止生成有害、偏见或不合规的内容。可以调用内容安全API或在后处理阶段进行过滤。工具调用沙箱化对于智能体调用的工具尤其是执行写操作、访问外部网络或执行代码的工具尽可能在沙箱环境中运行。例如使用Docker容器来隔离执行数据库查询的工具限制其网络和文件系统权限。3. 智能体层防护针对Prompt Injection与越权系统提示词加固在给AI模型的系统指令中明确、反复地强调其角色边界和安全规则。使用分隔符如|im_start|, 清晰划分系统指令、工具定义和用户输入。工具权限细分不是所有工具都对所有用户或所有会话开放。根据用户等级、会话类型动态地提供工具列表。例如普通用户不能使用“删除用户数据”工具。操作确认HITL如前所述对高风险操作强制引入人工审批环节。会话隔离确保不同用户的会话上下文绝对隔离防止通过提示词注入窃取他人会话信息。6.2 实战使用Redis实现分布式速率限制速率限制是防止API滥用的第一道防线。这里展示一个使用Redis实现的、可在分布式环境下工作的令牌桶算法。import redis import time class RateLimiter: def __init__(self, redis_client: redis.Redis, key_prefix: str rate_limit): self.redis redis_client self.key_prefix key_prefix def is_allowed(self, user_id: str, action: str, max_requests: int, window_seconds: int) - bool: 滑动窗口计数器算法。 user_id: 用户标识 action: 操作类型如 chat max_requests: 时间窗口内允许的最大请求数 window_seconds: 时间窗口大小秒 返回: True 如果允许请求否则 False key f{self.key_prefix}:{user_id}:{action} now int(time.time()) window_start now - window_seconds 1 # 使用Redis管道保证原子性 pipe self.redis.pipeline() # 1. 移除窗口之前的记录 pipe.zremrangebyscore(key, 0, window_start - 1) # 2. 获取当前窗口内的请求数 pipe.zcard(key) # 3. 如果未超限添加本次请求记录 current_count pipe.execute()[1] # 获取第二步的结果 if current_count max_requests: pipe.zadd(key, {str(now): now}) # 分数和成员都用时间戳 pipe.expire(key, window_seconds) # 设置Key的过期时间自动清理 pipe.execute() return True else: return False # 使用示例 redis_client redis.Redis(hostlocalhost, port6379, decode_responsesTrue) limiter RateLimiter(redis_client) user_id user_123 if not limiter.is_allowed(user_id, chat_completion, max_requests10, window_seconds60): raise HTTPException(status_code429, detail请求过于频繁请稍后再试。) # ... 继续处理聊天请求 ...这个实现使用了Redis的Sorted SetZSET。每个请求的时间戳作为成员和分数。每次检查时先清理掉当前时间窗口之前的记录然后统计剩余的数量。这种方法比简单的INCREXPIRE计数器更精确能平滑地处理时间窗口边缘的请求。6.3 持续的安全实践安全不是一劳永逸的配置而是一个持续的过程依赖项扫描使用safety、trivy或npm audit等工具定期扫描项目依赖修复已知漏洞。容器镜像安全如果使用Docker确保基础镜像来自可信源并保持更新。扫描镜像中的漏洞。凭证管理AI服务的API密钥、数据库密码等敏感信息必须使用环境变量或秘密管理服务如AWS Secrets Manager、HashiCorp Vault来管理绝不可硬编码在代码中。审计日志记录所有重要的用户操作和智能体工具调用包括谁、在什么时候、做了什么、结果如何。这些日志是事后调查和取证的唯一依据。这八天的集中攻坚从数据持久化到人类监督从实时推送到能力扩展再到全方位安全加固相当于为AI应用的后端搭建起了一个完整的“中枢神经系统”和“免疫系统”。每个环节都充满了细节和权衡但正是这些细节决定了应用最终的用户体验、稳定性和安全性。我的体会是在AI应用开发中“让功能跑起来”只是第一步而“让功能跑得稳、跑得安全、跑得可控”才是真正体现工程价值的地方。希望这份详细的复盘能帮助你在自己的项目中少走一些弯路。
返回列表