ARTICLE DETAIL

资讯详情

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

AI应用架构演进:从单体到事件驱动的实战指南

AI应用架构演进:从单体到事件驱动的实战指南 1. 这事得从一张“看起来很稳”的架构图说起先交代下背景。我这两年一直在帮中小自研团队做AI应用的技术落地自己也亲手搭过好几套系统。一开始大家的做法几乎一模一样——接到需求先上个单体服务把所有接口揉在一起用户请求进来代码里同步调用大模型API拿到结果直接返回前端。这套玩法在demo阶段简直爽得不行代码好写、调试直观、逻辑清晰本地一条命令就能跑起来。但一旦进入生产环境尤其是业务开始有点真实流量的时候问题就像开了闸一样往外冒。我自己印象最深的一次事故一个智能纪要功能用户上传一段两小时的会议录音前端跪等后端返回结果大模型调用超时HTTP连接直接断掉用户这边看到的是网络异常请重试。重试本身没毛病问题在于没有重试机制前端傻等、后端卡死、日志里全是半截调用记录。更尴尬的是当时系统里另一个模块的接口也共用了同一个线程池一个慢任务把整条请求链路堵得死死的别的正常功能也跟着遭殃。跟几个同行聊下来大家踩的坑几乎一模一样大模型调用慢、不稳定、偶发超时而单体架构把所有这些不确定性全部锁死在一条同步链路上一个环节抖动整条链路跟着抖。这时候大概都会冒出同一个念头——这架构是不是得动一动了。这篇文章想聊的就是我从实际项目里总结的一条演进路径从单体架构起步的AI应用如何一步步走向事件驱动架构。适合谁看中小自研团队的技术负责人、正准备给AI应用做架构升级的开发者以及那些想搞清楚事件驱动到底怎么落到AI应用里的人。我不打算堆概念尽量用我们真实踩坑和解决问题的过程来讲。2. 单体架构在AI应用里的真实处境2.1 为什么大家起步都喜欢搞单体这不是审美问题是生存问题。中小自研团队做AI应用通常是在验证一个业务假设这个智能体能不能解决用户的实际问题这个AI功能能不能带来留存或转化在这个阶段团队最稀缺的资源根本不是架构能力而是迭代速度。单体架构最大的优势就是心智负担低。一个服务搞定所有事代码都在一个工程里函数互相调用数据放同一个数据库debug的时候从上往下一路看过去就行。拿我做过的几个AI应用来说早期基本都是这种结构接入层收到用户请求解析参数业务层组装Prompt拼上下文调大模型API拿到结果后做后处理格式化、敏感词过滤、结构解析落库返回给前端整个过程就是一条直线出问题顺着调用链查就行。最快的时候一个功能从开发到上线只需要一天半。这个速度在早期很值钱因为它允许你用最低的成本快速验证这个AI功能到底有没有人用。说实话到现在我也依然认为早期AI应用就该从单体起步。一上来就搞微服务、事件驱动那不是设计那是给自己上刑。架构演进应该跟着业务复杂度走而不是跟着技术潮流走。2.2 当AI应用开始跑真实流量单体架构撑不住了但问题是业务一旦起来AI应用和传统Web应用的差异就会暴露得非常明显而这几个差异恰好都打在单体架构的命门上。第一个差异是慢。传统接口通常要求毫秒级响应大模型接口动不动就是几秒到几十秒。单体架构里的同步调用意味着每个请求进来一个线程就要被占住几十秒。假设你用默认线程池配置比如200个线程那基本上每秒只能处理个位数的并发请求。我见过最离谱的情况是一个只有50个用户的内部系统因为同事疯狂点按钮直接把后端线程池干满了。第二个差异是不稳。大模型API的延迟波动非常大高峰期可能比平时慢十倍。而且它偶尔会返回格式不对的内容——JSON多了一个逗号或是在一个应该输出纯文本的接口里突然吐出一大段Markdown。单体架构里应对这些不确定性的办法只能是加超时、加重试、加熔断但这些东西加在同步链路上会让代码越来越臃肿逻辑越来越绕。第三个差异是贵。大模型调用是按token计费的一次调用少则几毛多则十几块。单体架构里如果用户在页面等得不耐烦刷新了一下前端重新发请求后端就又调一次大模型。用户多刷几下账单就上去了。而且这种浪费是代码层面很难察觉的因为你看到的只是正常的请求日志根本意识不到里面有大量重复调用。这三个差异叠加在一起会让单体架构出现一个非常典型的现象单个功能可用性尚可毕竟每次调用都能成功返回但整个系统的吞吐量极低、成本失控、故障影响面不可控。说白了系统表面上活着实际上已经瘫痪了。2.3 AI应用架构演进的核心矛盾同步 vs 异步我自己的体会是AI应用架构演进本质上是处理同步思维到异步思维的转变。单体架构天然是同步的请求进来处理返回结果。这种模型适合短事务、快响应的业务比如登录、查询订单、获取用户信息。但AI应用的核心流程——调用大模型——天生是个慢操作。强行用同步模型承载慢操作结果就是线程被占住、链路被拖慢、用户体验差、系统资源被浪费。事件驱动架构恰恰是异步思维的产物请求进来先把任务变成一个事件发出去然后立刻告诉用户我收到了正在处理真正的大模型调用放到后台异步执行等结果出来后再把进度或结果推给用户。这件事想通之后架构改造的方向其实就非常清晰了把长耗时的AI调用从同步链路中拆出去用事件去驱动它。3. 事件驱动和AI应用为什么是天生一对3.1 先用人话讲清楚事件驱动大致是什么我给团队讲这个的时候从来不先摆概念都是讲点外卖。你点外卖的过程其实就是一个典型的事件驱动过程你下单事件产生商家接单消费者接收、备餐处理、骑手取餐配送继续处理而你不需要盯着商家炒菜你可以该干嘛干嘛等外卖到了再接收通知。中间任何一步做慢了你也不会一直卡在页面加载中的状态。事件驱动架构就是这么个逻辑。系统里各个模块之间不直接打电话等对方答复而是把要说的事写成一个事件丢到一个共享的地方消息队列/事件总线让感兴趣的人自己来取。生产事件的人和消费事件的人可以完全独立不需要知道对方的存在。放在传统Web架构里这相当于把同步调用改成了异步通知。放在AI应用里这个转变的价值就更大了因为AI应用里的菜做得特别慢——大模型推理不是几毫秒能搞定的事。3.2 AI应用里哪些环节该改成事件驱动根据我梳理过的多个AI应用项目下面这几类场景是最适合从同步调用改成事件驱动的。第一类是长文本生成类任务。比如AI写周报、AI生成营销文案、AI总结会议纪要。这类任务的共同特点是耗时从几秒到几十秒不等而且用户往往不需要实时等待结果做完通知一声就行。这类在单体架构里最痛苦因为一个长生成任务能卡死一堆线程。改成事件驱动后用户提交需求就立刻收到已受理响应后台异步跑大模型跑完把结果写入数据库、推送通知。第二类是流水线型AI任务。比如智能客服的完整处理流程意图识别 → 知识库检索 → 大模型生成答案 → 敏感信息检测 → 结果发送。每一步都可以单独做成一个事件消费者由事件总线把它们串起来。这样做的好处是每一步都可以独立扩缩容、独立重试。如果大模型生成这步因为API超时失败了只需要重发事件重新执行这一步前面的意图识别结果不用重新算。第三类是批量处理型AI任务。比如定时把一批用户消息做情绪分析、批量给商品生成标题和描述、后台把所有新上传文档做向量化。这类任务天然就是脱离用户请求的单体架构里只能用定时任务硬扛扛不住就是队列堆积、任务丢失。第四类是多智能体协作的编排任务。比如你搭了一个客服智能体它需要先判断问题归属再调用不同的子智能体处理最后汇总结果。这类任务里各智能体之间的通信本身就是天然的事件交互用事件驱动来编排比用代码硬编码控制流要灵活得多。上面这些场景基本就是当前AI应用开发里最主流的需求形态。如果你发现自己做的AI应用里有大量这类业务那就说明它的架构确实应该往事件驱动方向走一走了。3.3 事件驱动给AI应用带来的三个关键收益讲收益之前先泼个冷水事件驱动不是银弹它不会让你的接口变快甚至会让单个请求的端到端耗时变得更长。但它换来的是三个更重要的东西。第一个收益是削峰填谷。大模型API有它的吞吐上限而且经常在某个时段变慢。事件驱动架构里你可以把请求积压在队列里按大模型API能接受的速度去消费而不是让用户直接面对API的抖动。这就像水库上游洪水来了先蓄着下游慢慢放水系统不会因为瞬时流量就崩掉。第二个收益是故障隔离。单体的痛是一个功能挂全部功能挂。事件驱动架构里如果大模型服务挂了队列里的任务会堆积但其他功能比如登录、查询历史记录依然正常。等大模型恢复了队列继续消费任务也不丢。这在单体架构里是很难做到的。第三个收益是弹性扩缩容。事件消费者可以单独扩容。当队列积压严重时多起几个消费者实例吞吐马上上来闲的时候缩回去成本跟着降。这在单体架构里做不到——你不能只给AI调用部分扩容要扩就得整服务一起扩。4. 从单体到事件驱动我建议你这样动手改4.1 第一步先别动架构梳理你的调用链我见过很多团队一上来就引入Kafka、RabbitMQ恨不得把所有东西都解耦。结果呢业务没变复杂架构却先崩了。正确的第一步是梳理调用链。把你系统里所有调用大模型的地方找出来列个清单标注这几个信息这个调用是同步的还是可以有延迟的调用结果需要立刻返回给用户还是可以等完成后通知调用失败的影响是什么——是重试就能解决还是需要人工介入调用频率和耗时大概是多少做完这个梳理你就会发现AI应用里真正需要同步返回的调用其实很少。拿智能客服来说如果用户是在页面上实时对话那确实需要快速响应但如果是用户提交了一个投诉工单让AI帮忙生成处理建议那完全没必要同步等。我自己惯用的判断标准很简单用户的眼睛是不是盯着屏幕在等这个结果如果是保留同步调用如果不是就该改成异步事件。4.2 第二步选一个轻量的事件通道从小切口进去很多人一说事件驱动就想到Kafka我劝你先冷静。中小团队的AI应用初期流量大概率撑不起Kafka这种重组件运维成本会吃掉你的架构红利。如果技术栈是Java初期用RabbitMQ就挺好够轻、够稳、生态好如果在云上直接用云厂商的消息队列产品比如阿里云的RocketMQ、腾讯云的CMQ、AWS的SQSSNS更省心如果团队小、又不想引入外部中间件Redis Stream也是一个很不错的起步方案毕竟很多团队Redis已经有了不用新增基础设施。我第一次做改造时没有立刻引入独立消息中间件而是先在单体服务里加了一个事件中间层。这个中间层本质上就是个内存队列 一组Worker线程用Java里的ExecutorService加BlockingQueue就能实现。虽然这个方案不具备跨服务能力但帮团队把异步处理的流程先跑通了。等确认这套模式有效再接入RabbitMQ把内存队列替换掉。这种渐进式改造风险极小每一步都能单独验证、单独回滚。4.3 第三步用事件 状态机模式改造第一个AI功能我推荐第一个改造对象选择那些业务逻辑清晰、边界明确、失败影响可控的功能。比如我之前做过的AI周报生成就是个特别合适的练手对象。这个功能改造前是同步的用户提交一周的工作纪要后端同步调大模型生成周报前端转圈等待。改成事件驱动后整个流程变成了四个事件WeeklyReportSubmitted用户提交工作纪要接口层收到后生成一个事件发布出去立刻返回已受理周报生成中PromptContextPrepared消费者拿到原始纪要后组装Prompt补充角色设定、输出格式要求、历史风格参考组装完成发下一步事件LLMInvoked调用大模型生成周报这个消费者需要支持重试ReportPersisted把生成结果写入数据库然后通过WebSocket推送通知这块有一个很关键的细节每一步事件里都必须带上足够的上下文。因为生产者消费者解耦之后消费者手上只有事件数据没有共享内存可以读。你发布一个事件不能只发用户ID应该把用户ID、原始纪要、期望格式这些一次性打包在事件载荷里否则消费者还得回头去查数据库那解耦就白解了。状态机的作用是追踪这个任务走到哪一步了。我常用的做法是在数据库里建一张任务状态表记录任务ID、当前状态、事件历史、重试次数、下次重试时间。消费到事件后先查状态是否匹配匹配才继续往下走不匹配就丢弃——这是防止事件重复投递导致重复处理的防线。4.4 第四步把同步接口改造成提交-查询模式用户端体验也要跟着改。原来的同步接口是请求-响应改造成事件驱动后要变成提交-查询或提交-回调。具体做法是用户提交请求接口返回一个task_id前端拿到这个task_id后立即显示一个任务处理中的界面前端通过轮询查询接口每2-3秒一次来获取任务状态状态变成成功后前端再拉取结果。如果有WebSocket或者SSE条件可以让服务端主动推送状态变化体验更好。但对大部分内部工具型AI应用轮询已经够用了不用为了追求技术时髦给自己加复杂度。这里还有一个坑要提醒查询接口不要查两次数据库。我见过有团队查询任务状态时先查状态表再查结果表两个表数据不一致导致用户看到任务已完成但结果一直加载不出来。正确的做法是在状态表里直接存一个result_json字段状态变成成功时结果就已经在里面了查询接口一次搞定。5. 实操环节一个AI智能体应用的事件驱动改造实录5.1 改造前单体里一个典型的同步AI调用痛点说个咱们信息检索里频繁出现的场景——制度条例学习助手。这算是一个典型的AI工作室里搭建的智能体应用用户提问请帮我解释《项目管理制度》里关于变更审批的流程智能体需要先从知识库里检索相关条款再基于检索结果生成回答。单体架构的实现方式大致是这样app.post(/api/chat) async def chat(request: ChatRequest): # 1. 向量检索 docs vector_store.search(request.question, top_k5) # 2. 组装Prompt prompt build_prompt(request.question, docs) # 3. 调用大模型 answer llm.call(prompt) # 平均耗时5秒偶尔20秒 # 4. 返回结果 return {answer: answer}这个实现是不是看着很熟悉但跑生产环境之后这台制度条例学习助手会出现几个非常实际的问题用户在对话框里连发三个问题接口层就得同时占用三个线程去等大模型响应线程池很容易被打满大模型偶尔超时整个请求直接报500用户前面的问题也白问了不会自动重试没有上下文管理用户追问那审批需要几个工作日时系统不知道那指的是变更审批回答质量直线下降更重要的是这种设计无法应对知识库的更新。制度文件每天都在变向量库里存了十来个版本的制度文件没有版本筛选检索出来的内容新旧混杂AI给出的答案就是对错参半。这些都是不见到真实业务就根本意识不到的坑。5.2 改造后事件驱动的制度条例学习助手怎么工作我把这个应用做了一次事件驱动改造。改造后的结构是这样的用户提问进来接口层不再直接调大模型而是做四件事记录用户会话上下文生成一个question_id把用户问题 历史对话摘要 知识库版本号组装成一个QuestionSubmitted事件发布到事件通道立刻返回{question_id: ..., status: processing}后台跑三个消费者彼此之间完全解耦第一消费者专门做增强检索。收到事件后根据事件里携带的知识库版本号去向量库检索只返回该版本的制度条款。这一步做了多路召回向量检索用语义相似度找相关条款同时用BM25做关键词命中的补充召回把两路结果做合并排序保证短问句比如报销上限多少和长问句比如出差期间发生餐饮费用的报销标准是什么都能命中。然后它发布一个ContextReady事件事件里带着完整的上下文。第二消费者专门做大模型生成。收到ContextReady事件后把检索到的制度条款 用户问题 会话历史一起组装进Prompt调用大模型。这个消费者天然支持重试如果第一次调用超时事件走重试队列最多重试3次每次间隔递增。Prompt不是简单地把条款拼在问题后面而是明确让模型先判断检索内容里有没有权威答案如果没有就回答资料不足以回答该问题而不是硬编。这个设计能显著减少大模型一本正经地胡说八道。第三消费者专门做结果落库和通知。生成完成后把最终答案写入结果表并通过WebSocket推送回答已生成事件给前端前端拿到后主动拉取结果。5.3 关键代码事件通道和消费者的最小实现当时我不想在初期引入外部中间件用的Redis Stream做事件通道。生产端发布事件的代码大致是这样import redis, json r redis.Redis(hostlocalhost, port6379, decode_responsesTrue) EVENT_STREAM ai_events EVENT_GROUP knowledge_worker def publish_event(event_type: str, payload: dict): event { event_id: uuid.uuid4().hex, event_type: event_type, payload: payload, timestamp: int(time.time()) } r.xadd(EVENT_STREAM, {data: json.dumps(event)})消费端我写了一个通用的Worker基类核心逻辑就是轮询读取事件、按类型分发到对应的处理函数def start_consumer(consumer_name: str, handler_map: dict): while True: entries r.xreadgroup( groupnameEVENT_GROUP, consumernameconsumer_name, streams{EVENT_STREAM: }, count10, block5000 ) for stream, messages in entries: for msg_id, fields in messages: event json.loads(fields[data]) handler handler_map.get(event[event_type]) if handler: try: handler(event[payload]) r.xack(EVENT_STREAM, EVENT_GROUP, msg_id) except Exception as e: # 记录失败进入重试逻辑 log_error(event, msg_id, e)这段代码放到生产环境前有两个细节必须处理第一消费端要做幂等。事件通道在网络抖动时可能把同一个事件投递两次消费者收到重复事件时如果重复调用大模型钱就白烧了。我在事件处理前会先查一下任务状态表如果这个question_id已经处理过了直接ACK掉不再处理。第二消费端要做重试退避。事件处理失败不能立刻重试要按1秒、5秒、15秒、60秒的节奏递增。我直接给事件加了一个retry_count字段每次重试时更新它的下一次执行时间放到一个延迟队列里。不会用太复杂的方案简单可靠最重要。5.4 Prompt层面的配套改造异步化之后生成质量反而更高了这个点是我在实践里意外发现的。把所有环节异步化之后Prompt的构造余地反而大了。因为不再需要为了响应速度而压缩Prompt可以把知识库检索到的条款完整放进去不用只取前三条。制度条例助手这类应用对回答准确性要求高Prompt里多放几条候选条款做交叉验证生成结果的质量是要显著好于只给三条硬拼的。我在Prompt设计里加了两道防护一道是要求模型指出它依据的条例具体是第几条没有依据就明说没找到另一道是让模型将答案按结论-依据-操作建议三段来组织。这两道防护对学习助手、制度问答类应用极有用。因为这类场景用户拿着AI的回答做真实决策如果AI给了一条过时或者错误的制度描述那比不给答案还糟糕。异步化还让我有条件做流式交互之外的另一种体验用户提交问题后界面先展示状态轮播检索中 → 生成中 → 完成等完成后再展示完整答案。对于制度条例类问题这种模式反而比打字机式的token逐个蹦出来更符合用户预期——他们要的是一份能直接看的参考答案不是看直播生成过程。5.5 这套改造带来的实际效果改造完成后我做了几个对比测试数据分享出来给你们参考单次请求的接口响应时间从平均5秒大模型同步调用耗时降到了平均180毫秒只发个事件就返回用户端体感是提交后立刻得到响应大模型API的超时率从改造前的8%降到了1%以下因为超时之后会自动重试不需要用户手动再发一次系统支撑的并发问答量提升了将近一个量级因为线程不再被阻塞等待瓶颈从线程池大小变成了队列消费速度重复调用大模型产生的浪费大幅下降因为前端不需要在等待中反复刷新触发新的请求当然端到端的提问到拿到完整答案的时间其实差不多甚至略慢了一些因为多了事件传递和队列延迟。但用户体感反而更好因为马上有反馈和反馈后等待的体验远远好过长时间卡在加载中。6. 演进路上的坑和我的排查经验6.1 事件顺序问题制度条例两个版本答案打架刚做完异步化那阵子测试同学给我报了个bug用户先后问报销标准是多少和审批流程怎么走系统有时会把前一个问题生成的答案关联到后一个问题上答案错位。查了一圈发现是事件消费顺序的问题——同一个用户发起的多个事件在队列里不一定按顺序被消费。排查下来发现原因是我的消费者用了并发处理同一时刻pull了10个事件用线程池并行处理A事件的处理逻辑包含向量检索比较慢B事件的处理逻辑比较快先处理完了。结果B先把结果写进结果表A后写用户一查发现自己的第二个问题拿到了旧的数据。解决方案很简单但也很KISS按用户维度做事件分组同一用户ID的事件走同一个消费者实例且在该实例内串行处理。实现方法是在xreadgroup的时候指定消费者名字从user_id里算个hash保证同一个用户落到同一个消费者。这件事给我一个教训事件驱动灵活但顺序性这种天然需要强保证的东西你得在设计之初就明确要不要不要等出事了才想起它。6.2 事件无限重试钱是怎么烧没的还有个经典问题也是我踩得最疼的一次。某个配置错误导致大模型调用失败消费者一直在重试每次重试都打一次大模型API目的其实是为了看看是不是恢复但每次调用都是真金白银烧token。这个错误配置在公司里跑了一整夜第二天看账单的时候差点没背过气去。从那以后我做了三个硬性限制也写进了团队的SOP手册每个事件的重试次数硬编码为不超过3次超过就进死信队列可以便宜地人工处理每次重试的间隔用指数退避1分钟起步最多加到30分钟而不是隔几秒就疯狂重试重试前明确区分可重试错误超时、限流、临时网络抖动和不可重试错误Prompt格式错误、参数校验失败、事件数据缺失只有前者才走重试逻辑这里面的核心原则就一句话重试必须有上限、必须有退避、必须有原因判断。没有这三个约束的重试就是给云厂商送钱。6.3 事件积压监控看不到积压就等于没有积压事件驱动架构最需要盯的一个指标就是队列积压量。积压量大了说明消费者处理速度跟不上生产速度或者说消费者对应的下游服务比如大模型API出问题了。但这个问题在改造初期根本没引起我的注意因为队列偶尔积压稍微等一下就能消费完大家都没当回事。直到有一次大模型API持续故障了二十分钟队列里积压了上万条未处理事件等API恢复后积压事件疯狂消费直接把API的限流额度打穿了。那次之后我写了个监控脚本每分钟检查队列的pending数量超过阈值就告警。同时给消费者加了健康检查连续N次处理失败自动暂停新事件拉取避免一边消费一边失败一边无限循环。这里也建议你们省事的直接接云监控不用自己造轮子。6.4 别把全栈弄成全乱哪些环节保持同步反而是好事最后想提醒一句事件驱动不是要把系统里所有东西都变成异步。我自己试过把用户登录、权限校验也改成事件驱动结果给自己找了一堆麻烦。有些东西就天生适合同步。比如用户身份校验、核心链路的前提条件检查、必须实时返回给前端用于渲染的数据。这些环节做成异步反而会增加用户等待时间、把简单逻辑复杂化。我现在的原则是用户会直接感知等待的、失败后需要立刻反馈给用户操作的保持同步其他的大胆异步化。架构设计不是非黑即白单体里也可以有异步事件驱动里也可以有同步调用。好的架构师要会判断什么地方该用什么模式而不是拿着一把锤子看什么都是钉子。7. 说点我自己的体会做AI应用架构演进这几年最大的感悟是不要为了架构升级而升级。单体架构在项目初期完全够用因为它给了你极快的验证速度但当你的AI应用开始承担真实业务压力时——当大模型调用成了瓶颈、当用户开始因为等待而流失、当你发现失败重试和资源管理开始占据你的代码时——你需要的是架构上的松绑。事件驱动是这一步松绑里最实用的方向之一。演进没有那么玄乎。从梳理调用链、找出那些不该同步却同步了的AI调用开始引入一个轻量事件通道用一个功能做试点跑通流程后再逐步扩大范围。每个阶段都验证收益和代价每一环都能回退。我个人在实际项目中还有一个习惯每次架构改造完成一定要做一次回滚演练。事件驱动改造最怕的就是回不去了——万一事件通道挂了业务怎么办我在系统里留了一个开关可以临时把发事件切换回同步调用保证消息队列抖动时不至于让整个业务停摆。这个开关一开始只是应急方案后来反而成了稳定性兜底的一个关键设计。架构升级不是一锤子买卖你得给自己留退路。最后再分享一个小技巧给每个AI任务设计一个唯一的trace_id从事件发布开始经过检索、生成、落库每一步都用这个ID串起来打日志。排查问题的时候你会回来感谢这个决定的。分布式链路追踪在单体时代可有可无但一旦开始事件驱动它就不是可选项而是必需品了。
返回列表