ARTICLE DETAIL

资讯详情

深耕网站SEO优化与搜索引擎排名提升的一线实战洞察。

AI项目面试-高并发智能日志RAG系统

AI项目面试-高并发智能日志RAG系统 【项目背景】相关方店小秘主要做跨境电商ERP的4seller项目是欧美区的业务客户有几十万人用java写的。他们的日志系统分散在多个容器。跨服务调用比较难查。痛点开发人员排查效率低。跨部门沟通的时候研发人员被迫频繁中断手头的开发工作去充当“查单工具人”。同类异常再次发生时新员工依然需要从头摸索。目标客服输入自然语言系统直接回答。比如输入卖家 ID 或平台单号系统直接回答“该卖家今天上午有 50 个亚马逊订单拉取失败原因是 Amazon 接口限流”【确认数据来源数据口径定架构】1数据来源业务测的日志。2数据口径了解数据量大小了解业务流程关键节点可能会出问题的节点估算并发。依据我按业务流量放大来估算。比如订单核心日常峰值是 2000 TPS一个完整请求从拉单、校验、库存、支付到发货可能会产生 10 条日志那么日志 TPS 大概是2000 × 10 20000 条/秒。但跨境 ERP 还有回调、面单获取、库存同步、定时任务这些异步链路所以整体日志量会再放大。因此日常我们按 3万条/秒 做容量评估大促或第三方接口异常时按 10万条/秒级别做峰值预案。并发估算日常日志写入峰值1万 ~ 3万条/秒大促峰值5万 ~ 10万条/秒。查询接口并发10 ~ 50 QPSLLM 问答并发个位数到十几并发3设计架构【架构方案和技术选型】【redis】L1关键词缓存去重校验限流策略jwt 选型理由 高并发KV关键词检索简单高效。 【ES】L2语义缓存原始日志存储 选型理由 dense_vector 向量检索 keyword 精确过滤 range 时间过滤 bool 组合查询 【Milvus】向量数据库 选型理由 大规模日志向量检索 高并发向量召回 【Filebeat】监控日志文件读取新增日志发送到 Kafka 选型理由 轻量对业务侵入小支持断点续传直接感知变化 Filebeat 依靠 Registry 记录 Offset 保证“不漏”但受限于 ACK 机制和网络抖动必然会发生“重复发送”。企业级架构中不追求采集层的绝对精确而是通过在消费端生成唯一 log_id结合 Redis 或 ES 的 Upsert 做幂等去重来化解重复投递的问题。 Java 容器把日志目录挂载到了宿主机。 Filebeat 容器把宿主机的这个目录挂载进了自己的容器。 Filebeat 有读取权限。 【卡夫卡】解耦业务系统和日志系统抗住日志洪峰日志系统不会影响订单核心交易。 【高并发技术】 Web 查询层我们用 FastAPI 异步接口 日志写入链路用 Kafka 削峰 Python 用 asyncio 批量消费 写 ES 用 bulk 向量化部分用稳定 doc_id upsert防止重复生成。 【压测】 完整的全链路压测没有做得特别重但我们做过灰度验证。 先接入日志量较小的服务观察 Kafka Lag、ES 写入延迟、CPU 内存再逐步放大到核心订单服务。 ‌QPSQueries Per Second‌每秒查询/请求数统计服务器收到的独立接口请求数量通常不区分请求成败。 ‌TPSTransactions Per Second‌每秒事务数统计完整闭环业务成功完成的数量一个事务可能包含多个请求仅统计成功的事务。【离线阶段】【日志格式】{timestamp:2026-08-13T10:00:00.123,level:INFO,service:order-service,env:prod,trace_id:trace-abc123,seller_id:10086,shop_id:shop_001,platform:amazon,platform_order_no:111-1234567-1234567,internal_order_no:DXM_998877,event:order.pull.start,msg:开始拉取亚马逊订单} {timestamp:2026-08-13T10:00:01.456,level:INFO,service:order-service,env:prod,trace_id:trace-abc123,seller_id:10086,shop_id:shop_001,platform:amazon,platform_order_no:111-1234567-1234567,internal_order_no:DXM_998877,event:order.pull.success,msg:拉取亚马逊订单成功} {timestamp:2026-08-13T10:00:02.789,level:ERROR,service:order-service,env:prod,trace_id:trace-abc456,seller_id:20086,shop_id:shop_002,platform:shopee,platform_order_no:SP_556677,internal_order_no:DXM_778899,event:order.pull.failed,msg:拉取Shopee订单失败,error_code:REQUEST_TIMEOUT,exception:java.net.SocketTimeoutException: connect timeout}【数据样式】##1java传给我们数据的样子 { level: ERROR, service: order-service, env: prod, trace_id: trace-abc123, seller_id: 10086, platform: amazon, platform_order_no: 111-1234567-1234567, event: order.pull.failed, msg: 拉取订单失败, error_code: RequestThrottled, exception: AmazonApiException: rate limit exceeded, timestamp: 2026-08-13T10:05:23Z } ##2日志摘要我用摘要进行Embedding 【摘要格式】 服务: order-service平台: amazon事件: order.pull.failed 信息: 拉取订单失败错误码: RequestThrottled 异常: AmazonApiException: rate limit exceeded 卖家: 10086 【理由】 客服提问「张三的店铺怎么连不上了」 日志原文「AMAZON_API_AUTH_FAILED: token expired, store auth invalid」 ##向量化之后 { doc_id: hash(trace-abc123_order-service_1723543523), vector: [0.12, -0.34, 0.56, 0.78, ...], // 1024维 metadata: { trace_id: trace-abc123, seller_id: 10086, platform: amazon, event: order.pull.failed, error_code: RequestThrottled, summary: 服务: order-service平台: amazon事件: order.pull.failed..., timestamp: 2026-08-13T10:05:23Z } } ##知识库数据 { kb_id: kb_001, error_code: RequestThrottled, platform: amazon, category: 接口限流, root_cause: Amazon 平台对拉单接口有频率限制超出配额后返回限流错误。, solution: { steps: [ 1. 在配额平台查询当前卖家的剩余拉单配额。, 2. 启用指数退避重试策略参考配置 retry_config_v3。, 3. 若持续限流提交工单到平台对接组申请提升配额。 ], config_ref: retry_config_v3, escalation: 平台对接组 }, related_error_codes: [RequestThrottled, TooManyRequests], updated_at: 2026-08-01T00:00:00Z, kb_vector: [0.23, -0.11, 0.45, ...] // 用于语义召回 } 【用户问题】 卖家10086的Amazon店为什么拉单失败怎么解决 【相关日志证据】 服务: order-service平台: amazon事件: order.pull.failed 错误码: RequestThrottled异常: rate limit exceeded 卖家: 10086时间: 2026-08-13 10:05 【处理知识库经验】 错误码 RequestThrottled 处理 SOP 1. 查询配额平台确认剩余拉单配额 2. 启用指数退避重试配置 retry_config_v3 3. 持续限流则提工单到平台对接组 请基于以上信息给出排障结论和解决方案并引用来源。【ES细节】ES存原始日志和错误向量 order-log-* order-faild-log-*。手动按天存。容易漏查。性能差。使用Data Stream 取名order-logs创建 ILM 生命周期策略PUT _ilm/policy/order-log-policy单个索引保留半个月或者50G创建 Index Template匹配order-logs并开启data_stream。绑定刚才的 ILM 策略。Data Stream 必须有 timestamp。【python层】【输入】 { level: ERROR, service: order-service, env: prod, trace_id: trace-abc123, seller_id: 10086, platform: amazon, platform_order_no: 111-1234567-1234567, event: order.pull.failed, msg: 拉取订单失败, error_code: RequestThrottled, exception: AmazonApiException: rate limit exceeded } 【流程】 1读取 Kafka AIOKafkaConsumer 关闭自动提交 offset。 只有日志处理成功并写入 ES 后才手动 commit。 防止宕机丢日志。 2字段校验 如果日志缺少 trace_id、service、level 或 msg直接丢弃 3生成唯一 doc_id 基于 【trace_id 服务 时间戳 msg】 进行 MD5 哈希 4原始日志去重与批量写入 ES Data Streams 先用 Redis SETNX log:written:{doc_id} 过滤掉重复日志 然后放入 Bulk buffer前再批量写入 ES。 5判断是否需要向量化 生成摘要 只对 IMPORTANT_EVENTS (如 order.pull.failed) 等关键 ERROR 日志进行摘要转换 过滤掉无价值的 INFO 日志降低噪音。 摘要 【服务: order-service平台: amazon事件: order.pull.failed】 6Embedding 双 Key 去重与异步向量化 为了防止 Kafka 重复消费导致重复调用昂贵的 Embedding 模型采用 Redis 双 Key 机制 查完成 Key查 Redis processed:{doc_id}若存在说明已处理过直接跳过。 抢占处理 Key用 SETNX processing:{doc_id} (TTL 60s) 抢占分布式锁。若失败说明其他 Consumer 正在处理直接跳过。 异步向量化将摘要发给独立的 Embedding 服务/队列。 写入向量库并标记成功 向量写入 ES/Milvus 成功后 设置长 TTL 的 完成 Key (processed:{doc_id})。 释放锁删除 处理 Key (processing:{doc_id})。 7提交 Kafka offset 所有写入和去重逻辑执行完毕后手动 commit offset。 T0 Redis挂 → 连续3次失败 → 熔断器 OPEN T030s 冷却期到但还没人请求熔断器仍 OPEN T030sε 下一条Kafka消息进来 → _process_message → try_acquire → allow_request()发现冷却期已过 → 转 HALF_OPEN → 放行1个真实SETNX ├─ 成功 → 熔断器恢复 CLOSEDRedis确实好了 └─ 失败 → 立即回 OPEN再等30s consumer主循环每10s → _check_redis_replay() → 条件熔断器 CLOSEDis_healthy 队列有积压文件has_pending → 满足才调 replay_pending() → 逐条重新SETNX去重 → 补写ES【AIOKafkaConsumer理由】瓶颈根本不在 Kafka 的拉取速度而在后续写 ES 和调 Embedding 的网络等待时间。AIOKafkaConsumer 提供的协程异步 I/O 能力【ES单key机制去重】理由1由于 ES Data Stream模式无法指定自定义 _id我们要去重就要使用doc_Id去查询这样耗费性能2Kafka 多 Consumer 并发消费时的重复投递问题。度ES 原始日志写入Embedding 向量化日均数据量10 亿条1000 万条 (仅 ERROR)单次操作成本极低 (内部磁盘 I/O)极高 (GPU / API 费用)重复的代价极低 (多一条日志)极高 (浪费算力与金钱)防重策略单 Key 拦截(防并发)双 Key 防重(防重复计算)若加完成 KeyRedis 需数百 GB 内存 ❌Redis 仅需数百 MB 内存 ✅【Embedding-双key机制去重】并发场景10:00:00.000Kafka 突然抽风把同一条日志doc_idA1B2同时投递给了节点 1和节点 2。10:00:00.001节点 1查 Redisprocessed:A1B2不存在。于是开始调用 Embedding 模型耗时 500ms。10:00:00.002节点 2查 Redisprocessed:A1B2依然不存在因为节点 1 还在处理中还没成功没来得及标记完成。10:00:00.003节点 2认为这是一条新日志也开始调用 Embedding 模型。结果同一条日志被调用了两次 Embedding白白浪费了一倍的算力这就是并发穿透。宕机场景10:00:00节点 1抢到了锁processing:A1B2开始调用 Embedding。10:00:01突发状况节点 1 的 Python 进程突然 OOM 崩溃了或者服务器断电了。10:00:02锁没有被释放也没有任何标记说这条日志处理完了。结果如果这个锁有短 TTL比如 60 秒60 秒后锁过期了其他节点可以重试。但这依然有问题如果节点 1 其实没死只是网络卡了它最终处理成功了但锁已经过期被节点 2 抢走节点 2 又去处理了一遍还是可能重复。【在线阶段】用户提问改写后的 query │ ▼ ┌─────────────────────────────────────────────────┐ │ 【第一跳】检索原始日志 ──▶ 定位发生了什么 │ │ │ │ ① ES 结构化过滤 (seller_id / time_range) │ │ ② 混合召回: Dense向量 Sparse向量 BM25 │ │ ③ RRF 融合粗排 │ │ ④ Rerank 精排重排──▶ Top 5 日志 │ │ ⑤ 从日志中提取错误码: RequestThrottled │ └─────────────────────┬───────────────────────────┘ │ 拿到错误码 ▼ ┌─────────────────────────────────────────────────┐ │ 【第二跳】检索知识库 ──▶ 找怎么解决 │ │ │ │ 有错误码? │ │ ├─ 是 ──▶ 精确匹配知识库 (ES term / 映射表) │ │ │ 直接命中 SOP ──▶ 无需粗排重排 │ │ │ │ │ └─ 否 ──▶ 用 query 语义检索知识库 │ │ 召回多条 ──▶ (可选)重排 │ └─────────────────────┬───────────────────────────┘ ▼ ┌─────────────────────────────────────────────────┐ │ 【可选第三跳】检索历史案例库 ──▶ 上次怎么解决 │ │ 用错误码/向量检索相似案例 │ └─────────────────────┬───────────────────────────┘ ▼ ┌─────────────────────────────────────────────────┐ │ 组装日志(现象) 知识库(解决方案) 案例(参考) │ │ ──▶ 喂给 LLM ──▶ 完整排障结论 │ └─────────────────────────────────────────────────┘ { question: 最近 1 小时订单服务有没有报错, summary: 最近 1 小时 order-service 共发现 12 条 ERROR 日志主要集中在创建订单失败错误类型为 DatabaseTimeoutError。, logs: [ { timestamp: 2026-08-13T10:00:02, level: ERROR, service: order-service, trace_id: abc123, msg: 创建订单失败 } ] } 【为什么要Embedding错误日志】 很多错误日志根本没有 error_code 字段只有一段描述 日志「订单状态流转异常无法推进到已发货」 ← 没有错误码 日志「系统繁忙请稍后再试」 ← 没有错误码 假设遇到一个全新的错误码 PAY_ERR_999知识库里根本没有这条记录。 你的精确匹配链路精确匹配失败 → 返回暂无解决方案 → 系统显得很蠢。 向量检索把 PAY_ERR_999 的日志描述向量化去知识库找**语义最接近**的方案。比如匹配到了 PAY_ERR_100支付网关超时的处理方案作为参考。 客服问张三的店为什么出单这么慢 注意客服没给错误码。而你通过时间ID捞出来的3分钟日志里可能有十几个错误码 库存锁定失败 (INV_001) 物流下单超时 (LOG_005) 面单获取失败 (LBL_002) 订单状态异常 (ORD_009) 你的精确匹配链路不知道该匹配哪个错误码 → 懵了。 向量检索用出单慢去匹配发现物流下单超时LOG_005的日志向量最接近 → 精准定位根因是物流。 检索【历史相似案例】时 ⭐ 客服问之前有没有类似的出单慢问题 历史案例库通常是按问题现象组织的不是按错误码。精确匹配无能为力必须靠向量检索去找语义相似的历史案例。 【代码】 async def ask(question: str): # 1. 解析查询条件 query_params await llm_parse(question) # 2. L1 精确缓存 cache_key build_cache_key(query_params) l1_result await redis.get(cache_key) if l1_result: return l1_result # 3. L2 语义缓存 question_vector embedding(question) l2_result await es_search_semantic_cache(question_vector) if l2_result and can_use_cache(query_params, l2_result): # 回写 L1 await redis.set(cache_key, l2_result[answer], ex120) return l2_result[answer] # 4. 未命中走真实 RAG logs await search_es_logs(query_params) answer await llm_summarize(question, logs) # 5. 写缓存 await redis.set(cache_key, answer, ex120) await save_semantic_cache(question, question_vector, query_params, answer) return answer【记忆模块】【集群方面】【Redis 方面】我们用的是Redis Cluster 分片集群3 主 3 从。主要承载 L1 精确缓存、短期记忆和分布式锁。 在分布式锁的设计上我们特别注意了Hash Tag的使用比如emb:{doc_id}:processing保证相关 Key 落在同一个 Slot避免 Lua 脚本报 CROSSSLOT 错误。 因为短期记忆在 MySQL 有冷数据兜底所以即使 Redis 发生主从切换丢失少量热数据系统也能自动从 MySQL 降级预热业务完全可容忍。Redis Cluster 为什么是 3 主 3 从能不能 2 主 2 从“Redis Cluster 要求至少 3 个主节点才能完成 16384 个 Slot 的合理分配和故障自动转移。如果只有 2 主当其中一个主节点宕机时剩余的主节点无法形成多数派可能导致集群进入 FAIL 状态无法自动故障转移。3 主 3 从是生产环境的最小高可用配置。【ES 方面】我们部署了3 个专用 Master 节点 多个高配 Data 节点实现集群元数据与数据存储的职责分离。第一在节点角色上我们做了职责分离。部署了 3 个低配的专用 Master 节点只负责集群元数据管理和分片调度不存数据。这样即使 Data 节点因为复杂查询导致 CPU 飙升也不会影响 Master 的心跳彻底避免了集群‘脑裂’和错误的重平衡。第二在日志存储上我们采用了 Data Streams 配合 ILM索引生命周期管理。日志是典型的时间序列数据Data Streams 让应用层只需无脑追加写入底层自动按天滚动Backing Index。同时配合 ILM 策略实现热温冷数据自动分层最近 1 天的日志在 SSD 热节点7 天后自动迁移到 HDD 冷节点30 天自动删除。全程不需要写任何外部定时任务。第三也是最关键的我们将‘向量检索’和‘原始日志’做了物理隔离。因为向量检索KNN是极其消耗 CPU 的高维数学计算而原始日志写入和 BM25 查询主要消耗磁盘 I/O。如果放在同一个索引里AI 召回时的 CPU 峰值会直接拖垮日志的实时写入链路。所以我们单独建了向量索引实现了计算资源和存储资源的隔离保证了核心日志管道的绝对稳定。””服务器角色承载的分片说明Data 节点 1P1(主),R2(副本),R3(副本)存了分片1的主数据以及分片2、3的备份Data 节点 2P2(主),R1(副本),R3(副本)存了分片2的主数据以及分片1、3的备份Data 节点 3P3(主),R1(副本),R2(副本)存了分片3的主数据以及分片1、2的备份【集群概念】Master主节点 不需要存数据只是负责管理Data数据节点。数据节点里面有主分片和副本分布在不同数据节点 。写入原理分片和备份集群分布。用户随意的像一个节点发请求。之后这个节点就是协调节点。然后根据你的文档id。去hash一波。计算出你要去哪个分片。然后发给这个分片主分片。主分片处理完了。然后备份好。再响应给用户。查询原理分片和备份集群分布。用户随意的像一个节点发请求。之后这个节点就是协调节点。然后根据分发到所有分片。去查一波。协调节点汇总排序。发送get请求给这些分片要数据。要到了。再响应给用户。【如果 Data 节点 1 宕机了会发生什么】自动容灾机制会立刻启动第一步Master 节点感知故障。3 个 Master 节点发现 Data 节点 1 的心跳断了会在集群状态中把它标记为下线。第二步副本自动提升为主分片。第三步重建新副本。最终结果整个过程通常在几秒到几十秒内自动完成业务侧几乎无感知数据零丢失。”【系统评估】核心使用的是RAGAS 框架进行自动化评估。首先我们从历史真实故障中抽取了上百条案例构建了包含标准答案和关键日志上下文的黄金测试集Golden Dataset。然后每次调整 Prompt、更换 Embedding 模型或修改检索策略后我们都会跑一遍我们重点关注四个指标Context Recall上下文召回率确保 ES 和向量库没有漏掉那条最关键的 ERROR 日志。Faithfulness忠实度这是我们的上线红线指标。排障系统绝不能容忍幻觉我们要求 LLM 的结论必须 100% 基于召回的日志证据RAGAS 的 Faithfulness 分数必须大于 0.9 才允许发布。Context Precision上下文精确度监控召回的日志里是不是混入了太多无关的 INFO 噪音。Answer Relevance答案相关性确保 AI 直接回答了排障问题。对于 Query 理解层我们补充了自定义脚本评估槽位提取的 F1 分数对于工程侧我们使用压测工具监控P95 延迟和 Token 成本。通过这套‘RAGAS 自定义规则 性能压测’的组合拳我们保证了系统每次迭代的质量都是可量化、可追溯的。”Slot F1 低Query 理解层拉胯没提取出卖家 ID 或订单号。优化 Query 重写的 Prompt、增加正则表达式兜底规则、补充业务同义词词典。【web服务】┌──────────────────────────────────────────────┐│ 智能排障查询 ││ ││ 平台 [ 亚马逊 ▼ ] ← 下拉框 ││ 店铺 [ 10086-张三 ▼ ] ← 下拉框 ││ 时间段 [ 最近1小时 ▼ ] ← 下拉框 ││ 错误类型 [ 支付失败 ▼ ] ← 下拉框(可选) ││ ││ 问题描述 [ 付不了款____________ ] ← 自由文本 ││ ││ [ 开始排查 ] │└──────────────────────────────────────────────┘在线问答服务是基于 FastAPI 的 HTTP 服务。本地开发时我会直接用 Uvicorn 启动因为它支持异步 ASGI适合 FastAPI。但在生产环境我不会裸跑单进程 Uvicorn而是使用 Gunicorn 管理多个 UvicornWorker。这样一方面可以利用多核 CPU另一方面 Gunicorn 可以在 worker 异常退出后自动重启提高服务稳定性。通常启动命令是gunicorn app.main:app -k uvicorn.workers.UvicornWorker -w 4 --timeout 120。 因为我们后端还要调用 LLM 和 Rerank所以 timeout 会设置得比一般 Web 服务更长。离线 Kafka 消费服务不需要 Gunicorn 和 UvicornWorker因为它不是 HTTP 服务而是独立的后台 Consumer。离线服务会直接用 Python 进程运行或者通过 Kubernetes Deployment / Supervisor 来保证进程存活。SSE 适合 AI 文本流式输出它是基于 HTTP 的单向推送实现简单、网关友好、成本低WebSocket 适合实时语音、双向打断、多人协同和复杂任务推送。你们客服前端如果只是文本对话优先用 SSE等后续有实时双向需求再引入 WebSocket。1. WebSocket 连接管理2. 鉴权与 session 绑定3. 消息协议设计4. 心跳、断线重连、消息补偿5. 多实例广播 / 连接路由维度SSEWebSocket协议基于 HTTP独立协议HTTP Upgrade 后变成 ws/wss通信方向服务端单向推送到客户端客户端和服务端双向通信实现复杂度低高前端 APIfetch / EventSourceWebSocket API数据格式主要是文本 UTF-8文本、二进制都可以自动重连SSE 原生支持需要自己实现断点续传可用 Last-Event-ID需要自己设计鉴权fetch SSE 可以带 Header / Cookie浏览器原生 WS 不能自定义 Header通常用 query token网关兼容性更好本质是 HTTP需要支持 Upgrade部分代理/网关配置更麻烦负载均衡相对简单需要考虑长连接和会话保持适合场景AI 文本流式输出实时语音、双向交互、复杂状态推送运维成本低高企业级推荐文本对话首选实时语音/双向任务再引入
返回列表