一次查询,跨系统多模块数据并发聚合;规则引擎评估标签,AI Agent 生成可读结论;LLM 增强,规则兜底。我们用这套架构把单次工单处理从 2-3 分钟压缩到 20-40 秒,并让操作人员不再依赖「经验直觉」做判断。
导语
业务场景中存在这样一类典型难题:操作人员面对一个待处理工单,需要在 2 个系统间来回切换,逐一查看各自的 7-10 个模块,手动比对信息,再凭借经验综合判断,最终给出处理建议。熟练坐席处理一个工单通常需要 2-3 分钟,且结论高度依赖个人经验,不同操作人员给出的判断可能相差悬殊。
我们构建了一套基于 AI Agent 的智能诊断辅助决策系统,将这个流程重构为:一次请求触发跨系统多模块数据并发取数、规则引擎 + LLM 双通道标签评估、分层决策引擎输出带 Trace 的结论、AI Agent 生成自然语言总结与处理建议。各模块数据边算边推、实时呈现,坐席无需等待全部完成即可开始分析,整体处理时长压缩到 20-40 秒,关键路径提效约 3-9×。此外,SSE 流式推送让坐席可以在等待期间并行处理多个工单,进一步摊薄单工单的实际时间成本。
本文从架构设计、关键工程决策、踩坑与复盘三个角度展开,聚焦五个核心设计模块:配置驱动的取数编排、规则 + LLM 双通道标签评估、分层决策引擎(带 Trace)、SSE 实时事件流、AI Agent 总结生成。
由客服场景核对确认用户身份信息,然后审核决策,总结落地的一种 Agent 诊断推荐决策方案,可类比多接口数据查询分析对比决策场景。
1.1 典型场景的三类痛点
我们要解决的是「多维度信息聚合 + 综合研判」场景,这类场景在客服、风控、审核等领域普遍存在。以操作人员处理一个标准申请流程为例,痛点集中在三个方面:

三个问题叠加,导致数据分散在各处却难以协同利用——即便每个系统都在正常运转,整体处理效率依然低下。
💡 这不是数据不够,而是缺一套把数据聚合、理解、决策串联起来的引擎。

1.2 量化损耗

2.1 挑战一:并发取数与依赖调度
不同数据源之间存在隐性依赖:A 接口的返回结果是 B 接口的入参,B 未就绪时 C 无法启动。但全部串行调用代价太高,等待时间是各接口耗时之和。
单纯「全并发」也行不通——部分接口需要上游返回的字段作为参数,并发启动会拿到空参数,导致取数失败或污染缓存。
2.2 挑战二:判断逻辑的结构化表达
经验型判断难以代码化有两个原因:
规则本身是模糊的——「绑定时间够不够长」「异地登录算不算风险」,这些判断人脑里有隐含的上下文,翻译成代码容易出现边界问题。
数据缺失与条件不满足是两件事——传统
if-else无法区分「没有这个数据(无法评估)」和「有数据但不满足条件(已评估为不满足)」,两者混为一谈会导致误判。
2.3 挑战三:诊断结论的可读性与可追溯性
机器给出的结论如果只是「通过」或「拒绝」,操作人员无法理解为什么,更无法在结论有争议时提供依据。要让机器的判断「可信赖」,必须让每一步决策都有 Trace,且这个 Trace 必须是人类可读的。
💡 可解释不是锦上添花,而是人机协作的基础。没有可读 Trace,操作人员永远不会真正信任机器的判断。
3.1 核心设计原则
经验规则化、数据可追溯、决策可解释。
三个目标贯穿整个设计:
把人工研判经验结构化沉淀为规则引擎,而不是锁在个人经验里;
多维度数据聚合后形成可追溯的判断依据,每条标签都能说明「基于哪个数据得出」;
每步决策都带 Trace,操作人员能逐层看到「为什么走到这里」,建立可信赖的人机协同。
3.2 系统总体架构

3.3 核心设计分层矩阵

4.1 配置驱动的取数编排:Domain as Config
🎯 挑战:跨 2 个系统、共 14-20 个模块的数据单元,各有接口参数、依赖关系、缓存策略、字段处理规则。硬编码每一个不可维护;新增一个数据维度需要改代码、发版。
💡 解决方案:把「一次底层数据查询」抽象为一个独立的取数单元(Domain),单元的所有行为完全由配置驱动。
🛠 实现细节:
每个 Domain 的配置包含 7 个维度:

配置片段示例:
{ "key": "unit_profile", "biz_label": "用户画像", "command": "", "params": { "target_id": "$req.target_id" }, "cache_ttl_sec": 172800, "needs_upstream": [], "on_success": { "extract": { "profile_field_a": "$resp.data.field_a", "profile_field_b": "$resp.data.field_b" }, "triggers": ["unit_risk_check"], "post_handler": "PostProfileTranslate" }, "cache_check": { "field": "code", "allow_values": ["0"] }, "mask_fields": [ { "field": "profile_field_a", "perm_key": "", "mask_fn": "cover_middle" } ]}并发编排的核心代码(orchestrator.go):
// 一次性并发启动所有 Domain,每个 Domain 内部处理自己的依赖顺序var wg sync.WaitGroupwg.Add(len(scheduleKeys))for _, key := range scheduleKeys {key := keygo func() {defer wg.Done()result, err := EnsureDomain(orchCtx, key, taskID, req, registry)// ...取数完成即推送 SSE 事件sendEvent(orchCtx, taskID, "detail_ready", map[string]any{"domain": key,"data": normalizeResult(result),})// ...域完成即评估仅依赖本域的标签_, cards := EvaluateTagsByDomains(orchCtx, req, partialData, []string{key})for _, c := range cards {sendEvent(orchCtx, taskID, "tag", c)}}()}wg.Wait()
📊 效果:新增一个数据维度只需增加一行配置,零代码改动;编排器通用化,统一支持所有数据单元的并发调度。
💡 配置即代码,声明即执行。新增数据维度的成本从「改代码+发版」降为「改 JSON 配置文件」。
4.2 进程内并发合并 + 跨实例分布式锁:双层去重
🎯 挑战:高并发场景下,同一请求的同一 Domain 可能被多个 goroutine 并发触发(取数扇出、后置处理、断线重连等场景均可能重复触发)。重复取数不仅浪费下游资源,还可能把竞态结果写入缓存。
💡 解决方案:两层去重,由近到远:

进程内并发合并的关键实现(ensure_domain.go):
// 每个「请求ID + Domain」对应一个 domainState,通过 sync.Map 保证唯一stateI, _ := stateTable.LoadOrStore(sk, &domainState{done: make(chan struct{})})st := stateI.(*domainState)if atomic.CompareAndSwapInt32(&st.status, 0, 1) {// 抢占成功:执行真实取数result, fetchErr := executeDomainFetch(ctx, key, taskID, req, cfg, reg)if fetchErr != nil {stateTable.Delete(sk) // 失败立即清理,允许下次重试}close(st.done) // 通知所有等待者} else {// 等待他人完成select {case// 复用结果casereturn nil, ctx.Err()}}
一个关键细节:取数失败时立即删除 stateTable 中的记录,允许下次重试;取数成功则延迟 30 秒清理(窗口内的重复请求都能复用内存结果,无需再走 Redis)。
📊 效果:相同请求内的 Domain 重复触发率降为 0;Redis 读写次数减少约 60-80%。
4.3 缓存可信度校验:错误响应不入缓存
🎯 挑战:下游原子接口在业务异常时会返回非零业务码(如「对象不存在」「权限不足」等),如果把这类响应写入缓存,后续请求会直接拿到错误结果,且在 TTL 内无法纠正。
💡 解决方案:每个 Domain 可配置 cache_check,仅当响应中指定字段命中白名单值时才写缓存:
// isCacheable:白名单校验,不在列表内则跳过缓存func isCacheable(ctx context.Context, cfg *entity.DomainConfig, result any) bool { if cfg.CacheCheck == nil { return true // 无配置:始终允许缓存 } m := toAnyMap(result) raw, ok := m[cfg.CacheCheck.Field] if !ok { return true // 字段不存在:保守允许 } for _, allow := range cfg.CacheCheck.AllowValues { if toStr(raw) == allow { return true } } log.WarnContextf(ctx, "[EnsureDomain] 业务码校验不通过,跳过缓存 域=%s 值=%v", cfg.Key, raw) return false}💡 白名单永远比黑名单可靠。「允许缓存哪些状态」比「禁止缓存哪些状态」更安全——新增业务码时不用同步维护黑名单。
4.4 标签评估引擎:规则 + LLM 大模型双通道
🎯 挑战:有些标签判断逻辑清晰(绑定时长、数值比较),可以用规则表达;有些标签语义模糊(登录地点是否「相关」),规则难以覆盖。同时,数据缺失时传统的二态(命中/未命中)会导致误判。
💡 解决方案:规则通道处理结构化判断,LLM 通道处理语义模糊判断;两者都输出三种状态:正向命中、负向命中、数据不可用。
两个通道如何分工:

三种结果状态:

这个设计解决了一个实际的工程问题:数据缺失 ≠ 条件不满足。
不区分缺失状态的写法:
// 危险:location_match 为 false 可能是「确认不匹配」也可能是「LLM 未返回/接口超时」if !h["location_match"] { // 误判:把「没有数据」当成「地点不匹配」}显式区分缺失状态的写法:
// locationState: "pass" | "fail" | "unknown"// 仅当明确 pass 时才视为通过,unknown 保守处理locationState := detectLocationState(h)locationOK := locationState == "pass"
📊 效果:误判率(因数据缺失导致的错误决策)降为 0;LLM 标签的 Prompt 可通过远程配置中心热更新,无需发版。
4.5 分层决策引擎:P0-P4 优先级 + 全流程 Trace
🎯 挑战:研判逻辑不是扁平的「满足 N 个条件即通过」,而是有明确的优先级——某些条件一票否决,某些条件可以弥补其他条件的不足。需要表达这种优先级逻辑,同时让每步判断可追溯。
💡 解决方案:五层优先决策引擎,命中即返回,不继续下探;每层都记录 DecisionStep。

可解释 Trace 的实现:决策引擎不只返回最终结论,还返回每一步的判定明细:
type DecisionStep struct { Level string // "P0" / "P1" / ... Name string // "P1-强关联通过" Passed bool // 是否命中该层 Reason string // 一句话判定原因 Facts []string // 关键事实(对象特征、状态等) HitTags []string // 命中的标签(中文标题) MissingTags []string // 缺失的标签(如有)}Trace 有两个用途:
前端直接展示:操作人员能逐层看到「为什么走到 / 为什么没走到这一层」,建立信任感;
作为 AI Agent 的输入上下文:让大模型以完整的判定过程为依据,生成有针对性的自然语言解释,而非空泛结论。
📊 效果:决策逻辑完全显式化,新增规则只需修改决策引擎对应层次的判断逻辑,影响范围清晰;Trace 日志同时服务于运营审计和问题排查。
4.6 SSE 事件流 + 断线续读:实时感与可靠性兼得
🎯 挑战:诊断过程涉及跨系统多模块并发取数 + LLM 总结,整体耗时 20-40 秒。如果同步等待所有数据就绪再返回,坐席体验极差;但如果只做简单的 SSE 推送,客户端断线重连后会丢失已推送的事件。
💡 解决方案:进程内事件总线(实时推送)+ Redis 事件流(持久化回放)双轨并行。

断线续读的实现要点:
// writeSSEEvent:每条事件携带全局自增序号fmt.Fprintf(w, "id: %d\nevent: %s\ndata: %s\n\n", seq, eventType, data)// 重连时:客户端通过 Last-Event-ID 告知已收到的序号// 服务端先从 Redis 补齐历史事件,再接上实时流// 同时做序号去重,保证「不丢不重」
订阅必须早于编排器启动:如果先启动编排器再建立订阅,编排器产出的早期事件会在订阅就绪前被错过;先订阅再启动编排器,则不存在这个问题。
幂等处理:
// 同一请求的重复调用按任务状态机处理switch status {case Running: // 不重复启动,仅订阅已有事件流case Done: // 清空旧事件流后重跑(命中缓存,秒级返回)case Pending: // 正常启动编排器}📊 效果:前端「边算边看」体验;断线重连后无缝续接,无事件丢失;重复请求幂等处理,不产生额外的下游接口压力。
💡 「边算边推」解决实时感,「持久化回放」解决可靠性,两者是同一问题的两个维度,缺一不可。
4.7 AI Agent 总结生成:为什么不直接调 LLM
🎯 挑战:诊断流程最终需要把标签评估结果、决策 Trace 汇总成坐席可直接使用的自然语言建议。直接拼 Prompt 裸调 LLM 在工程上存在几个明显问题:模型配置硬编码难以热更新、多个调用场景(标签评估 + 最终总结)各自维护连接浪费资源、流式输出需要自己处理事件循环、Prompt 和模型参数调优需要发版。
💡 解决方案:引入 Agent 框架,把「模型管理 + 会话管理 + 流式执行」的工程复杂度下沉到框架层,业务层只关注输入什么、期望输出什么。
Agent 框架带来的具体工程价值:

共享 Runner 的实现(agent_service.go):
// 服务启动时创建一次,全局复用——标签 LLM 评估与最终总结共用同一 Runnerfunc NewAgentService() *AgentService { sessionService := inmemory.NewSessionService() // 内存 Session,天然隔离不同工单 agt := agents.CreateAgent(ctx) agentRunner := runner.NewRunner(appName, agt, runner.WithSessionService(sessionService)) // ...}按调用场景动态覆盖 Prompt(summary.go):
// Prompt 优先从远程配置中心读取,调整话术无需发版prompt, _ := rainbowConf.GetRainbowKv(ctx, group, "summary_prompt")agentOpts := []agent.RunOption{ agent.WithAppName("qqzm_summary"), agent.WithInstruction(prompt), // 覆盖默认 instruction,每次调用独立}eventChan, _ := req.Runner.Run(ctx, sessionUser, sessionID, message, agentOpts...)流式输出 + 失败降级:Agent 以事件流方式逐步输出内容,同时通过 SSE 实时推送给前端;解析失败时自动降级为规则引擎结论,核心流程不中断:
// 流式逐 chunk 回调,同步累积最终答案_, _, err := agents.ProcessStreamingResponseWithCallback(ctx, eventChan, func(delta string) { finalAnswer.WriteString(delta)})// 解析失败 / 字段为空 → 降级为规则引擎结论if err != nil || json.Unmarshal([]byte(raw), &out) != nil { return fallbackSummary(matchedRule), nil}输出清洗:模型有时把「1. xxx 2. xxx」粘连成一段,用正则拆成独立行,坐席直接可读:
// 兼容「句号+空格」「句号无空格」「仅空格」三种粘连形式var numberedSuggestionRe = regexp.MustCompile( `(?:\A|\s+|([。.;;,,!!??))]))\s*(\d+)\.\s*`)
💡 Agent 框架解决的不是「让 LLM 更聪明」,而是「让 LLM 调用工程化」——把模型管理、会话隔离、流式执行、Prompt 热更新这些胶水代码下沉到框架,业务层只写业务。
4.8 双层权限模型:可用性与合规性的显式取舍
🎯 挑战:权限校验失败有两种情况性质完全不同——登录态失效(该拦截)vs 权限数据源故障(不该因此把所有人挡在门外)。字段级权限则相反,需要严格不降级。
💡 解决方案:两层权限,两种降级策略显式对立:

字段掩码引擎在接口出口统一执行,敏感字段明文仅在缓存中保存,下发前端前按权限决定是否掩码:

💡 降级是一种设计决策,而非被动容错。显式声明「哪里降级、哪里不降级」,比隐式兜底要安全得多。

关键日志示例:
[EnsureDomain] Redis缓存命中 域=unit_profile 耗时=2ms[EnsureDomain] 业务码校验不通过,跳过缓存 域=unit_risk 值=1001[Orchestrator] 决策完成 level=P1 conclusion=强关联证据充分...[Summary] 最终总结生成完成 耗时=3.2s
6.1 量化提效

⚡ 量化提效:单次工单处理从 2-3 分钟压缩到 20-40 秒,提效 3-9×;SSE 流式推送支持坐席边看边分析、并行处理多工单,进一步放大实际吞吐效率;维度扩展和规则调优从「改代码发版」变为「改配置即生效」。
6.2 能力下沉

📥 能力下沉:资深人员的经验从「人脑私有」变为「规则公共资产」;新人上手周期大幅缩短;系统维护的部分工作从开发前移到运营侧。
6.3 模式升级

📚 模式升级:从「经验黑箱 + 手工查询」升级为「规则透明 + 并发聚合 + 可解释决策」——这不是效率的优化,而是工作方式的重构。
这套系统的工程设计有五个值得在类似场景复用的准则:
7.1 配置即代码:Config as Code
取数单元、标签规则、Prompt 均以声明式配置存储,逻辑由通用引擎解释执行。新增/修改业务维度不触碰运行时代码,降低变更风险,缩短迭代周期。
适用场景:规则频繁变化、维度持续扩展、非技术人员需要参与规则维护的系统。
7.2 数据缺失不等于不满足:Unknown as First-Class Citizen
任何依赖外部数据的判断,都应把「数据拿不到」作为独立状态显式标记,而不是当成「条件不满足」处理。少了这一状态,接口超时或数据源故障时系统会静默误判。
适用场景:依赖多个外部数据源、数据可用性不保证、误判代价高的判断系统。
7.3 决策即 Trace:Decision as Trace
决策引擎的输出不是一个结果,而是一条证据链。每步判定记录「命中了什么」「缺失了什么」「为什么」,Trace 既直接展示给操作人员,也作为 AI Agent 生成自然语言总结的输入依据,同时支持运营审计。
适用场景:需要解释决策原因、支持人工复核、机器与人协作的系统。
7.4 降级即设计:Fallback as Design
区分「该降级的地方」和「不该降级的地方」,并在设计阶段显式声明,而不是靠运行时兜底。基础可用性层降级,合规红线层严格不降级。
适用场景:可用性与合规性存在张力、权限体系复杂的系统。
7.5 实时即持久:Realtime as Persistent
长耗时异步流程中,实时推送(进程内事件总线)和持久化回放(外部存储事件流)必须并行,缺一则要么丢事件、要么断线后无法恢复。订阅必须先于编排器启动,避免早期事件被遗漏。
适用场景:耗时较长(> 5 秒)、需要流式展示进度、客户端网络不稳定的场景。
适合本方案的场景
数据分布在多个系统、多个模块,且模块间存在依赖关系
判断逻辑可以拆解为「优先级规则 + 少量语义模糊标签」
操作人员需要理解决策原因,而非只看结论
维度和规则需要频繁迭代,不希望每次都发版
不适合本方案的场景
纯 LLM 端到端场景(原因:此方案的核心价值是规则可解释,若规则本身就无法结构化,配置驱动的引擎意义不大)
实时性要求极高(< 100ms)的场景(原因:并发取数 + 标签评估 + LLM 总结的链路有固有延迟,不适合低延迟 SLA 场景)
完全自动化执行(无人工确认)的高风险操作(原因:系统定位是辅助决策而非替代决策,不可逆操作不应由系统自动执行)
💡 核心判断:如果你的场景可以用「多源数据 + 规则研判 + 人工确认」描述,这套方案值得参考;如果是「端到端自动化」或「纯实时决策」,请寻找更适合的架构。
9.1 数据维度的持续扩展
当前架构对新增数据维度的扩展成本已经很低(配置文件驱动),下一步可以探索:
支持动态配置热加载(当前需要重启服务才能生效),目标是配置修改后秒级生效;
增加跨数据单元的聚合计算能力,支持多个模块的数据交叉计算。
9.2 决策引擎的自学习
当前决策规则完全靠人工定义。长期看,可以:
收集操作人员的人工复核结果作为标注数据;
用历史数据分析规则阈值的合理性,提供调优建议(而非自动调优——规则变更仍需人工审核)。
9.3 主动触发而非被动等待
当前系统是「有申请才启动诊断」的被动模式。未来可以探索在特定条件满足时提前预取、预缓存,进一步压缩操作人员等待时间。
这套系统的核心不是技术的堆砌,而是一次工作方式的重构:把散落在多个系统和多个人脑中的知识,通过配置化的规则引擎、可解释的决策 Trace 和 LLM 辅助,重新组织成一个可验证、可追溯、可复用的整体。
我们在这个过程中最深的体会,是「可解释性」的价值——不只是让操作人员理解机器在做什么,更是让机器的判断真正成为操作人员工作的一部分,而不是一个他们不得不执行却不明所以的黑盒指令。
机器给出建议,人负责决策,专业的交给专业的来,模型天然适合做各类信息整理,合并聚类等重复机械式任务,这些都交给模型处理完,最终由人来决策应该怎么处理,怎么安抚用户等等,各司其职才是最和谐的人机协同方式。
五条工程准则(配置即代码 数据缺失单独标记 决策留完整记录 降级策略显式声明 实时推送与持久回放并行)可独立拆用,不必整套照搬——哪个准则解决了你当前最痛的问题,就先从那里入手。
