
先交代一下背景。我所在的团队之前做过一版基于规则引擎的客服机器人后来业务量上来之后用户问题越来越发散规则维护成本直线飙升逼着我们把技术栈换成了LangChain 大模型。系统上线后原来的单机同步调用开始扛不住流量接连出现请求超时、进程假死、用户排队等待时间过长这些问题。后来我们花了几周时间专门针对高并发场景做了流控、排队和语义降级这三块能力的建设总算把系统稳定在了可用状态。这篇文章就是把我们在生产环境踩过的坑、最终落地的方案和关键代码整理出来希望能给正在用LangChain做大模型客服或者类似场景的朋友一些参考。1. 整体架构设计思路先解决“并发进来后会发生什么”在设计一个高并发智能客服系统之前最先要想清楚的不是“用哪个模型”也不是“怎么调LangChain”而是“当100个、500个甚至1000个用户请求同时涌进来的时候我们的系统会经历什么”。1.1 问题的本质大模型接口不是免费的无限资源大模型服务的调用成本很高这个成本体现在两个维度第一是费用维度。调用一次商业大模型API是按token计费的并且多轮对话中每轮都要把历史消息拼接进上下文这会造成token消耗成倍放大。一个10轮对话的session可能单次请求就要消耗2000到4000个token这比单轮问答贵得多。第二是性能维度。大模型的推理速度天然比传统接口慢。一个常规的生成式回答往往需要1到5秒才能完整返回遇到长文本生成甚至能到10秒以上。如果系统不做保护任由并发请求直接穿透到模型服务层那么很快就会出现两类问题模型服务端主动限流抛HTTP 429或者503错误导致大量请求失败更可怕的是慢请求堆积请求在应用层排队等待模型响应一个接口的延迟会被放大到所有调用方。所以我们在设计之初就已经明确了一个基本原则所有用户请求不直接触达大模型链路必须先经过一道可控的流量入口由入口统一规划请求何时能进入模型服务、应该进入哪个降级策略。1.2 技术选型为什么最终选了异步消息队列而不是纯同步限流第一版时我们用的是最简单的方案——在API网关层用令牌桶限流超过阈值的请求直接返回“系统繁忙”。实测下来业务方完全不能接受原因很简单用户问一句“我的订单什么时候发货”系统回一句“系统繁忙”这个体验太糟糕了。更严重的是如果用户连续重试重试请求会加剧系统压力形成恶性循环。这里的关键认知是限流不是目的让有限的计算资源服务好更多请求才是目的。如果直接拒绝超出阈值的请求那只是把压力还给了用户并没有真正提升系统的吞吐能力。所以后来我们调整为“异步化 排队”的思路。请求进来后先进入Kafka消息队列由专门的工作线程按一定速率消费再调用LangChain链路完成模型交互最后通过回调或轮询把结果返回给用户。这个方案的好处有三个削峰填谷。业务流量天然是不均匀的618大促或者产品突发bug时流量可能是平时的10倍。队列能把瞬时高峰填平让后端模型服务始终工作在一个稳定的速率区间而不是忽高忽低。请求不丢失。队列模式下即使某个后端实例崩溃消息还在队列里其他实例可以接续消费这对客服场景非常重要——用户的问题不能丢。便于实现排队状态展示。我们可以把排队的position、预计等待时间这些信息透传给前端用户体验比“系统繁忙”好得多。架构上我们把系统拆成了接入层、队列层、消费处理层、模型层和结果回写层。一句话概括就是外面是流量入口中间是一个蓄水池后面是稳定的出水口蓄水池的出水速度是可以精确控制的。1.3 整个系统的分层结构为了让后续的内容更容易理解我先把最终的分层结构在这里交代清楚。第一层是接入层。也就是我们的API服务用户通过Web、小程序或者App发来问题。接入层只负责做基础校验、写入TraceID用于全链路追踪、把消息投递到Kafka然后立刻返回一个任务ID给前端。这个阶段的耗时应该控制在10毫秒以内绝不能在这里同步等模型出结果。第二层是排队层。排队层实现的是两件事一个是对消息做优先级分类比如VIP用户、投诉类消息可以插队另一个是对每个用户做幂等处理同一个用户短时间内的重复提问要合并或去重。第三层是消费处理层。这是LangChain主战场多个消费者实例从Kafka拉取消息每条消息经过前置意图识别、槽位填充、多轮上下文管理、RAG检索等流程最终组装prompt并发给模型。第四层是模型接入层。这层需要适配不同模型服务。我们实际生产环境同时接了私有化部署的Qwen和外部商业模型API不同模型有不同的限流配额这一层要做的就是针对每个模型做独立的流控。模型返回结果之后再经过结果解析、敏感信息校验、格式化通过WebSocket或者消息推送给前端。整个过程是异步的但用户感知上其实和同步问答差距不大。2. 流控策略设计控制流量进入模型的节奏2.1 流控的目标是保护下游不是用户体验的敌人很多初学LangChain的朋友容易有一个误区认为流控就是限制QPS把一切可能超过阈值的请求拦掉。实际上在生产环境中流控的核心目标是保护下游服务高可用同时保证系统整体的吞吐能力被善用。我们系统要流控的对象其实是两个一个是LangChain处理链路的消费速率另一个是对模型服务本身的调用速率。消费速率控制的是“我们以多快的速度接收用户消息”模型调用速率控制的是“我们以多快的速度把prompt发给模型”。这两者必须配合。如果消息消费得快但模型调用被限流那系统内部就会出现积压如果消息消费得慢后端的算力又闲着浪费。所以我们在设计上让这两层用两套不同的令牌桶参数并且允许动态调整。2.2 令牌桶算法的落地实现我们最终采用的是经典的令牌桶算法来实现流控。令牌桶的核心思想是系统以一个固定的速率往桶里放令牌桶的容量有限每个请求到来时必须从桶里取到一个令牌才能继续执行如果桶空了请求要么等待要么被拒绝。相对于漏桶算法的均匀流出令牌桶允许一定的突发流量因为桶里可以攒下最多capacity个令牌。这个特性对客服场景很有用——突然涌入的一小波流量可以被桶内的存量令牌消化掉不至于一有波动就触发限流。下面这段代码是从我们的流控组件里精简出来的实现了核心的令牌桶逻辑import threading import time class TokenBucket: 令牌桶限流器 rate: 每秒生成的令牌数 capacity: 桶的容量即最大突发流量 def __init__(self, rate, capacity): self.rate rate self.capacity capacity self.tokens capacity self.lock threading.Lock() self.last_refill_time time.time() def _refill(self): now time.time() delta now - self.last_refill_time self.tokens min(self.capacity, self.tokens delta * self.rate) self.last_refill_time now def acquire(self, blockTrue, timeoutNone): 尝试获取一个令牌 blockTrue时会阻塞直到获取到令牌 blockFalse时获取不到就直接返回False if not block: with self.lock: self._refill() if self.tokens 1: self.tokens - 1 return True return False start_time time.time() while True: with self.lock: self._refill() if self.tokens 1: self.tokens - 1 return True if timeout is not None and time.time() - start_time timeout: return False time.sleep(0.01)这段代码有几个点我想说明一下锁是必要的。我们的消费处理是并发执行的如果不加锁多线程同时修改tokens会出现竞态条件导致限流失效。_refill方法采用惰性计算只在每次获取令牌时才去计算这段时间应该生成多少令牌不需要后台线程去定时加令牌节省了资源。blockFalse的方式适合用在接入层做快速判断blockTrue适合用在消费处理层做等待让消费线程在令牌不足时短暂休眠而不是直接丢弃请求。以我们生产环境的经验为例接入层配置的流控参数是rate20即每秒钟最多有20个请求能通过接入层的初始校验。capacity10表示允许最多10个请求的瞬间突发。这里capacity10不是随便写的。我们结合业务场景做过分析以每天早上10点整的场景为例可能会有密集的新会话创建但更多情况是用户看完上一条回复后的思考时间普遍在1.5秒以上所以瞬时突发在10个以内是可以接受的。超过10个的部分就会进入排队逻辑而不是直接拒绝。2.3 动态流控根据模型响应时间自适应调整固定速率的限流在流量模型稳定的场景下效果不错但大模型服务的响应时间波动很大——同样的请求可能模型高峰期要6秒才返回低峰期2秒就返回了。如果消费速率是固定的很可能出现这种情况模型已经慢了我们还在用同样的速率往里发请求导致模型服务端队列暴涨。于是我们后来又加了一层动态流控根据最近一段时间内模型调用的平均响应时间和错误率动态调节令牌桶的参数。实现的思路简洁描述如下实时统计最近1分钟内模型调用的成功率success_rate和平均响应时间avg_latency我们设定了两个保护阈值比如成功率低于95%或平均响应时间超过5秒则进入自我保护模式在自我保护模式下把消费处理层的消费速率降为原来的50%同时缩小桶容量每5分钟检查一次如果指标恢复正常再逐步恢复消费速率——每次恢复增加10%直到回到原设定值这样缓启动可以避免一恢复就把系统再次压垮。这里需要分享一个我们踩坑的细节限流参数不是拍脑袋定的要考虑业务用户能接受的等待时间。比如一个用户发消息后如果5秒内没收到“对方正在输入”或者排队提示大概率会退出页面或者重复发送消息。所以在设计流控时要通过业务数据反推用户能接受的最长等待时间是多少、平均每轮问答模型耗时多少这两个数值之差就是消息在队列里允许滞留的时间上限。基于这个上限再结合业务峰值QPS才能算出合理的排队容量预留。另外我强烈建议流控组件不能做成进程内的单体组件。我们第一版就是把TokenBucket写成了单机实例部署两台服务器之后发现总QPS是单机限流的2倍一度被自己的限流策略坑了。后来我们把流控数据改存到Redis中用Lua脚本保证原子性实现在分布式环境下的统一限流。3. 排队机制实现既不让请求丢失也不让用户干等3.1 队列设计为什么用Kafka而不是Redis List我们最初同时在测Kafka和Redis List两种消息队列方案最后选择了Kafka。如果只是简单的“先进先出”消息堆积Redis的List结构其实也够用但客服系统有几个Kafka能更好满足的需求第一是消息回溯。Kafka支持基于offset重新消费万一消费程序出现bug或者LangChain链路上游模型服务大面积故障我们可以把队列里的消息重新消费不会导致用户消息凭空丢失。Redis List里的消息一旦被LPOP弹出去如果没有同步做备份是真的找不回来了。第二是多消费者组。我们的业务除了实时客服外还有离线分析场景。同一个用户消息既需要被实时链路消费来生成回答又需要被数据链路消费去做意图统计和服务质量分析。Kafka的消费者组机制天然支持这些一份数据可以被多个业务独立消费互不干扰。第三是分区有序。Kafka可以在生产消息时指定key比如用用户的user_id作为key这样同一个用户的消息会进入同一个分区消费时就能保证同一个用户的消息是按顺序处理的。这个特性在多轮对话的客服场景中很重要——两个消费者的线程如果同时处理同一个用户的两个消息上下文就乱了。在创建Kafka主题时我们做了如下配置# 创建topic10个分区3份副本 kafka-topics.sh --create --bootstrap-server kafka-01:9092 \ --replication-factor 3 \ --partitions 10 \ --topic customer_service_requests \ --config cleanup.policydelete \ --config retention.ms86400000分区的数量10是这样确定的我们的消费端预期最多部署5个实例每个实例的并发消费线程数设为2那么最多同时有10个消费线程在跑。分区数大于或等于最大并发消费线程数才能保证每个分区都被有效消费。分区数太少会导致部分消费线程空闲分区数太多会加剧Kafka的rebalance开销10是一个在测试后选择的平衡值。retention.ms86400000表示消息保存1天。客服场景一般没有必要超过这个时长因为用户如果真的等了超过1分钟还没收到回复我们后续还会通过离线工单体系兜底。3.2 优先级与插队策略VIP用户不能永远排最后做过客服系统的都知道如果队列严格按先来后到排队VIP用户投诉“我已经等了半小时”时会非常尴尬。所以我们的排队逻辑必然要支持优先级。我们的设计是在消息体中加一个priority_level字段取值范围是0到10默认是5。VIP用户的消息进来时Level直接提升到9投诉类的消息通过意图识别预判提升到8。实现“插队”这件事并不能直接作用于Kafka——Kafka的消息是分区有序的不支持在Topic级别做优先级重排。我们的方案是在消费处理层自己做一个本地优先级缓冲队列。具体做法是每个消费线程从Kafka拉取一批消息下来先不直接处理而是放到一个按优先级排序的本地队列中然后每次从本地队列头部取一条消息来执行LangChain链路。这样Kafka只负责消息的可靠传输和削峰真正的优先级调度放在内存里完成毫秒级延迟完全可控。本地优先级队列的代码简化后如下import heapq import itertools class PriorityMessageQueue: def __init__(self): self.heap [] self.counter itertools.count() def put(self, priority, message): # 使用heapq实现最小堆priority值越小越先被处理 # 我们外部已经把VIP的priority换算成更小的值 entry [priority, next(self.counter), message] heapq.heappush(self.heap, entry) def get(self): if not self.heap: return None # 弹出优先级最高的消息 _, _, message heapq.heappop(self.heap) return message需要强调的是这里之所以要用counter做次级排序是为了防止两个消息优先级相同时出现比较错误。Python的heapq比较元组时会依次比较每个元素如果前两个元素都相同就没有可比较对象了加上一个严格递增的序号可以保证任何两个元组都能比较出大小。初版我们没加这个counter结果在Python 3下直接抛了TypeError后来才补上的。3.3 排队状态透传让用户知道自己在第几位如果你做过客服或者工单系统就知道用户最怕的不是等待而是“不知道要等多久”。所以我们单独设计了一套排队状态查询机制。用户请求进来后接入层会立刻生成一个task_id返回给客户端。客户端拿到这个task_id后通过WebSocket连接订阅任务状态变更。处理链路的中间状态会实时写入Redis缓存包括WAITING_IN_QUEUE: 消息还没有被消费者拉取PROCESSING: 正在进行意图识别或模型调用SUCCEEDED: 回答生成完成FAILED: 链路失败进入了降级流程。在WAITING_IN_QUEUE状态时我们还会异步计算该任务在队列中的大概位置和预计等待时间。计算公式是预计等待时间 当前队列积压消息数 / 当前消费速率 历史平均单条处理耗时这个值不需要很精确给用户的感觉是“我们关注你并且知道你在等待”。实测这个机制上线后客服渠道的用户投诉率降低了近30%因为很多用户看到“前方还有3人在排队预计等待20秒”的时候是愿意等下去的——他们会理解系统正在为他们服务。3.4 超时控制与死亡消息处理排队最怕的就是消息进入队列后消费端出了问题一直消费不动结果用户在另一端等到天荒地老。所以消息的超时控制必须处理好。我们给每条消息设置了从投递到开始被消费的有效时间默认是60秒。这个值比业务上可接受的最大延迟要大一些但小于用户登录态会话的失效时间。如果消息在本地缓冲队列里超过10秒未被处理这里我们做了多级检查就触发超时监控报警并且把该用户的会话标记为“需要人工介入”。消费端处理消息时LangChain链路设置了整体超时30秒——这个30秒是多个环节的预算总和。其中意图识别调小模型耗时约2秒知识库检索向量数据库查询约500毫秒大模型推理端到端约8到15秒输出解析和内容安全校验约1秒。如果总耗时超过30秒消息会被判定为超时并投递到专门的timeout_ordersTopic中。有一个独立的补偿任务会扫描这个Topic将超时原因记录到监控系统同时触发前面的降级策略比如给用户推荐一个FAQ列表而不是一直干等。很多朋友问我们为什么要搞这么复杂直接同步请求、设置一个30秒的HTTP超时然后返回结果不是更简单吗但问题在于一旦业务量上来同步模型会让服务器的进程数和线程数线性膨胀每个线程都占着内存去等模型响应最终系统会陷入“线程爆满→内存升高→GC频繁→请求更慢→更多线程堆积”的死循环。异步队列从本质上杜绝了这种情况因为处理线程数是可控的、恒定的。4. 语义降级当模型不可用时如何让客服“还能用”4.1 降级不是“对不起我累了”而是有策略的预案在设计客服系统时一个很容易被忽略但又最影响口碑的点是当链路出现故障时系统对用户表现出来的行为是什么。很多系统的降级策略就是简单粗暴地返回一句“系统繁忙请稍后再试”。这句话对用户来说基本意味着“你们不行”。而我们在实践中总结出的经验是即使在故障期也必须有“努力为用户解决问题”的姿态。这个目标的实现方式就是语义降级。所谓语义降级指的是同一个用户问题的处理链路可以按多级别的计算资源来执行每一级能提供的回答质量不同。当高一级的资源不足时系统自动降级至低一级而不是直接失败。我们在系统里设计了五级降级路径完整大模型链路 RAG检索最理想的处理方式能深度理解用户意图结合内部知识库生成详细回答。简化大模型链路无RAG如果知识库检索服务向量数据库出问题就走这一步仅靠大模型通用能力回答不读取内部知识。小模型意图分类 高频问答库匹配如果大模型API不可用就把问题交给一个轻量级的文本分类模型先识别出用户问题所属的标准问题ID然后从FAQ库中直接命中标准答案。这块在我们系统中由Embedding相似度计算加倒排索引来完成。兜底关键词规则匹配如果向量检索也出了问题就退化为传统的关键词匹配通过同义词扩展匹配标准问题。人工客服工单转接以上层级全部失效时把用户问题转为工单记录完整的上下文承诺“客服会在XX时间内回复”。在LangChain里这种多级降级的实现思路可以借助LangChain的RunnableBranch或者自定义BaseCallbackHandler来设计流程。我后面会给出实际代码。4.2 流程编排在LangChain链路中嵌入降级判断我们在生产环境中的LangChain代码逻辑实际上不是一条线性链而更像是基于决策的流程图。为了便于阅读我在这里给出一段简化版的代码展示了如何在一个方法中按顺序尝试不同层级的解决方案并且通过异常捕获和状态检查实现自动降级import time from langchain.prompts import ChatPromptTemplate from langchain.schema.runnable import RunnableLambda from langchain_community.chat_models import ChatOpenAI from langchain_community.embeddings import OpenAIEmbeddings from langchain_community.vectorstores import FAISS class SemanticDegradeChatService: 语义降级客服服务 def __init__(self): # 初始化不同级别的客户端 self.llm ChatOpenAI( modelqwen-plus, temperature0.3, max_tokens512, timeout20, # 这里用openai兼容的base_url接入其他模型服务 base_urlhttp://your-model-gateway.internal/v1, ) try: self.vector_store FAISS.load_local( knowledge_index, OpenAIEmbeddings(modeltext-embedding-v2), allow_dangerous_deserializationTrue ) self.vector_available True except Exception: # 向量库不可用标记降级 self.vector_available False self.vector_store None # FAQ库question - answer模拟高频问答 self.faq_data { 退款到账时间: 退款通常在3-5个工作日内原路退回。, 如何修改收货地址: 请您在订单发货前在订单详情页点击修改地址。, 发票如何获取: 您可以在订单完成后在财务中心申请电子发票。, } self.keyword_map { 退款: 退款到账时间, 发票: 发票如何获取, 地址: 如何修改收货地址, } def _vector_rag_answer(self, question): 第一级RAG完整链路 if not self.vector_available: raise RuntimeError(vector store unavailable) docs self.vector_store.similarity_search_with_score(question, k3) if not docs or docs[0][1] 0.6: # 相似度分数不达标视为找不到知识 raise RuntimeError(no relevant document found) context \n.join([doc.page_content for doc, score in docs]) prompt ChatPromptTemplate.from_messages([ (system, 你是智能客服助手请仅基于以下背景资料回答问题\n{context}), (human, {question}) ]) chain prompt | self.llm return chain.invoke({context: context, question: question}).content def _llm_only_answer(self, question): 第二级纯大模型回答不检索知识库 prompt ChatPromptTemplate.from_messages([ (system, 你是智能客服助手请用简洁友好的语言回答用户问题。), (human, {question}) ]) chain prompt | self.llm return chain.invoke({question: question}).content def _faq_match_answer(self, question): 第三级FAQ语义匹配 # 实际项目中这里会调轻量级向量模型或者ES的more_like_this查询 # 简化版直接走关键词初步匹配再计算jaccard相似度 best_std_question None best_score 0.0 for std_question, _ in self.faq_data.items(): score self._cosine_similarity_text(question, std_question) if score best_score: best_score score best_std_question std_question if best_std_question and best_score 0.4: return self.faq_data[best_std_question] raise RuntimeError(faq not matched) def _keyword_match_answer(self, question): 第四级关键词规则匹配 for keyword in self.keyword_map: if keyword in question: std_answer self.faq_data[self.keyword_map[keyword]] return f为您找到相关解答{std_answer}如未解决请转人工 raise RuntimeError(keyword not matched) def answer(self, question): 提供回答内部按级别尝试 if not question: return 请描述您的问题 degrade_levels [ (RAG完整链路, self._vector_rag_answer), (纯LLM回答, self._llm_only_answer), (FAQ匹配, self._faq_match_answer), (关键词匹配, self._keyword_match_answer), ] last_exception None for level_name, func in degrade_levels: start_time time.time() try: result func(question) elapsed time.time() - start_time # 记录降级等级与耗时 print(f[Answer] level{level_name}, elapse{elapsed:.2f}s) return result except Exception as e: # 每一级失败都记录日志便于排查 print(f[Degrade] level{level_name} failed: {e}) last_exception e continue # 全部失败最后转人工工单 self._create_manual_work_order(question) return 非常抱歉暂时无法自动回答您的问题。 \ 我们已经生成工单客服将在1小时内回复。 staticmethod def _cosine_similarity_text(text1, text2): # 简化版文本相似度 # 生产环境建议用tfidf向量或者模型embedding set1 set(text1) set2 set(text2) if not set1 and not set2: return 1.0 intersection len(set1 set2) union len(set1 | set2) return intersection / union if union 0 else 0.0 def _create_manual_work_order(self, question): # 写入工单系统... pass在流控和排队已经确保整体稳定的基础上语义降级就是应对局部故障的另一个安全网。answer()方法从最高价值的RAG链路开始尝试每失败一次就跳到下一级在每一级都做异常捕获和日志记录保证不管哪一层出了问题用户那边总能得到一个“还在努力解答”的正向反馈而不是冰冷的错误提示。这个简化的实现里_cosine_similarity_text用的是字符交集计算只是一个demo。生产环境中我们用的是轻量级Sentence Embedding模型输出128维向量向量存储用的是Redis Search单次查询耗时约20毫秒左右即使每天百万级请求也不会有压力。4.3 降级状态的可观测性别等到用户投诉了才发现降级策略本身不难实现难的是知道“什么时候系统正在降级”。如果降级发生了但监控面板上没有信号那这个降级机制就是不完整的。我们给每个降级级别定义了明确的Metrics指标cs_answer_total: 所有回答数量标签区分levelrag/llm/faq/keyword/manualcs_answer_degraded_total: 发生降级的次数cs_answer_elapsed_seconds: 各级别耗时直方图cs_upstream_errors_total: 模型或知识库上游错误数。每次降级跳级时都会给Prometheus Counter打点并配套设置告警规则。比如当FAQ匹配路径的调用占比连续5分钟超过60%时就自动触发告警通知到值班群。这个告警告诉我们虽然系统还在工作但用户已经享受不到我们最好的回答质量了这往往是上游模型服务已经异常的信号。另一个需要特别留意的是降级策略不能无限递归。比如在FAQ匹配失败后跳到关键词匹配如果关键词匹配又调用了大模型来做同义词扩展那就是变相把降级路径又引回到故障的模型服务上。我们专门检查过所有降级路径确保每一级使用的技术栈是向下隔离的如果第1级依赖外部大模型API第3级决不能依赖同款API。4.4 限流与降级的配合以“请求预算”的角度统一规划最后一块让我想明确的是流控、排队和语义降级不是三个独立的功能模块它们在设计中应该是彼此配合的整体。我习惯用一个“请求预算”的概念来理解它们之间的关系。大模型服务链路就像一条金贵的生产线每个请求要消耗一定的计算预算。流控负责决定预算的分配上限——比如每秒最多花多少个请求的预算排队负责处理预算暂时不足时的人员调度而语义降级则是在预算彻底不够时换个更低成本的生产方式去交付成果。在我们线上系统中这三者的配合机制如下请求进入接入层先经过分布式令牌桶做总量限制。此时系统会记住请求的degrade_available标记如果当前链路健康度评分较高允许走完整链路如果模型服务的健康分已经低于阈值则直接跳过链路入口直接进入FAQ匹配层。如果请求通过了限制进入Kafka队列排队。消费端根据自己的处理速度拉取消息一旦发现上游模型服务的成功率低于阈值或P99延迟超过预设值消费线程会启动自我保护降低消费速率。消费线程拿到消息执行LangChain链路。链路内部按照前面介绍的五级路径进行尝试每级失败则降级到下一级。用这样的组合拳我们的客服链路在单机消费能力不高的前提下整体支撑住了日均百万级的请求量极大缓解了高并发场景下的基础设施压力。前面所有设计的价值和目的都统一到了这个体系里面让每一个有效请求都能在用户可接受的等待时间内获得一个可用的回答做到关键时刻客服“不掉线”这就是整个工程的最高目标。5. 实战中的常见问题与排查技巧最后一部分我把我们上线大半年以来实际踩过的一些典型问题和排查思路整理出来。这些问题在官方文档和各类教程里很少被提及但它们往往才是真正决定系统能不能稳定运行的关键。5.1 Kafka消费者堆积但CPU和内存利用率却很低现象监控面板显示Kafka消费延迟持续增长消息堆积数从几百涨到了几万。但一看服务器CPU利用率不到10%内存也很正常。排查过程起初我们怀疑是消费者线程数不足于是把单机消费者线程从2个调到了10个结果堆积依然存在CPU还是上不去。后来定位到问题出在流控组件的blockTrue逻辑上。我们当时配置了消费端在获取令牌时使用阻塞模式而令牌桶的rate只有5。也就是说整个系统每秒最多处理5条消息但上游QPS却有20Kafka自然越积越多。CPU空闲是因为大部分线程都在time.sleep(0.01)等待令牌。解决方案这是一个典型的流控参数和业务吞吐预期不匹配的问题。我们把消费速率从5逐步上调到30并同步评估了模型服务的P99延迟和错误率确认上调后模型服务依然稳定。又加了动态流控机制避免模型变慢时消费速率依然过快。经验排查队列堆积问题时不能只看消费者的处理速度还要看限制消费者处理速度的上游令牌桶参数。同时要把流控参数和业务目标挂钩不能用一套固定值应对所有流量变化。5.2 LangChain多轮对话中上下文错乱现象用户连续问几个问题比如先问“订单多久发货”再问“那退款呢”系统给出的回答和第一轮毫无关联。进一步排查发现不同HTTP请求被负载均衡分发到了不同的后端实例而LangChain的多轮对话上下文保存在进程内存中导致A实例不知道用户在B实例问过什么。解决方案我们统一在接层根据用户身份和会话ID从Redis读取历史消息组装成LangChain所需的messages数组。每轮对话结束后把最新的历史对话重新写入Redis并设置过期时间为30分钟。这样同一个用户不管请求落在哪个实例上都能取得完整的上下文。伪代码思路如下def build_messages(session_id, user_input): history redis_client.lrange(fchat_history:{session_id}, 0, 9) messages [] for msg in history: messages.append(HumanMessage(contentmsg[human])) messages.append(AIMessage(contentmsg[assistant])) messages.append(HumanMessage(contentuser_input)) return messages这里有一个值得注意的地方LangChain中对多轮对话有不同的存储方式可以持久化到多种存储介质我们实际上使用的是存储消息聊天的专用接口但核心要点是一致的——历史记录必须放在进程外部的共享存储中否则分布式部署时上下文就飘忽不定。5.3 大模型API偶发超时导致整条链路被拖垮现象商业模型API偶尔会出现单次请求超过20秒甚至30秒的情况。我们一开始给LangChain LLM设置了30秒超时但一个慢请求会占住一个消费线程一旦出现多个慢请求消费线程池就被全部占满后续消息全部积压。解决方案我们做了三层保护在调用模型API前先向Redis申请一个“并发占位”用信号量控制同时最多只有5个请求在模型上执行超出的请求直接在本地排队在LangChain LLM客户端设置更严格的超时初始为15秒如果模型网关反馈有熔断迹象动态下调为8秒增加“半开”重试机制单次请求失败后不立即重试而是连续失败N次后熔断一段时间等模型服务恢复后再放量重试。代码上通过LangChain的动态回调机制实现超时检测其中一个方式是from langchain.callbacks.base import BaseCallbackHandler class TimeoutCallback(BaseCallbackHandler): def on_llm_start(self, serialized, prompts, **kwargs): self.start_time time.time() def on_llm_end(self, response, **kwargs): elapsed time.time() - self.start_time if elapsed 10: print(fWARNING slow llm call: {elapsed:.2f}s)排查到这里建议大家务必记住一个原则在生产系统中对模型API的超时必须比单纯匹配HTTP超时更严格否则你的系统会以很隐蔽的方式慢死。而且模型调用的并发数需要独立控制不能用消费线程数作为模型并发数。5.4 用户重复点击导致重复消费与重复回答现象前端用户发了一条消息因为没马上收到回复又点了几次发送结果Kafka队列里同一个用户短时间内出现了多条几乎一样的消息。系统处理完成后用户收到了多条一样或不完全一样的回复。解决方案在消息生产端做去重。每条消息带一个由前端生成的client_msg_id接入层在写入Kafka之前先检查Redis中该用户最近2秒内是否有相同的client_msg_id如果有则直接丢弃并返回原任务ID。同时在消费端也做一层幂等处理前将task_id client_msg_id写入Redis设置过期时间为10分钟如果处理中发生重复消息Redis中的SETNX命令会保证只有第一个请求能抢到处理权。# 消费端幂等 def consume_message(message): dedup_key fdedup:msg:{message[client_msg_id]} acquired redis_client.set(dedup_key, 1, nxTrue, ex600) if not acquired: return # 已经处理过或正在处理 process_message(message)这一层非常重要因为大模型问答比普通接口更耗时用户在等待期间的重复点击概率非常高不去重的话会导致无效的模型调用既浪费钱又污染上下文。5.5 语义降级后用户问题匹配率过低需要尽快补量与调优现象一次大模型服务故障持续了20分钟我们触发了FAQ匹配和关键词匹配的降级路径。故障恢复后看数据发现降级期间FAQ匹配的命中率只有30%出头大量问题最终转到了人工工单。原因分析高频FAQ库的覆盖度不够。很多用户问题说法很口语化比如“我的货到哪了”“什么时候给我发货”但在FAQ库里对应的标准问法是“查看订单物流信息”两者表面差距很大因此匹配效果很差。解决方案我们做了两件事。第一把历史客服会话中的高频用户原话提取出来人工标注后挂到相应的标准问题下扩充FAQ的同义表达集。第二把关键词匹配从简单的字符匹配换成基于词向量的相似度匹配先用轻量级embedding模型把用户问题编码成向量再和FAQ库中的向量做ANN检索实测在降级场景下命中率提升到了60%以上。想提醒大家的是降级方案需要提前准备和持续运营不能等到故障发生了再去扩充FAQ库。我们在每次大促前都会做一次降级演练用高仿真压测流量模拟上游模型故障观察降级路径下的各项SLO借机补全FAQ和关键词库。这些工作平时看不出价值但真正出故障时就能看出差距了。6. 写在最后几点工程化心得做这个高并发智能客服系统的过程中我最大的感受是LangChain确实大幅度降低了搭建大模型应用的编码和抽象成本但真正决定一个客服系统能否在生产环境扛住流量、稳定运行的关键其实在LangChain之外——在于流控、排队、降级这些偏传统后端工程能力的打法上。要说最核心的一条经验就是别把大模型服务当成一个能无限扩容的普通HTTP接口它更像是有产能上限的稀缺生产线。围绕这个基本假设把流量管理的思路和语义降级的预案前置到系统架构设计阶段会比事后打补丁高效得多。另外还有一点想分享每一次降级路径的触发都是产品体验最直观的试金石。我们把降级视为一种功能在日常工作中反复演练并尽可能优化它。最终目标不是让降级不发生而是让降级发生时用户依然能得到一个最接近正常水平的回答。我们这套系统的后续规划中还在探索把“排队等待时间预测”和“用户答题意图预判”结合在排队期间就预生成一些可能的回答选项进一步缩短用户的感知等待时间。如果后续有成果我会再把这部分经验整理出来分享。