# Risk 架构总览
信息源:Risk 业务架构图与请求生命周期简要介绍(2026-07)& Risk 知识库
目录
0. 先装 AI 问答(推荐)
这份文档只讲骨架。细节、代码、排障,直接问 AI —— clone 风控 AI 知识底座 risk_af_agent_hub,装到你的 Cursor / Codex:
git clone ssh://gitlab@git.garena.com:2222/shopee/loan-service/credit_backend/risk_af_agent_hub.git
cd risk_af_agent_hub
bash install.sh --agent codex ~/你的项目目录
bash setup-env.sh # 首次按引导配 Token装完直接问,例如:
- credit_risk_engine 的工作流怎么执行的
- risk_feature 的取数链路,从配置到下游走一遍
- 贴一段报错日志让它定位
内置知识库 + risk 服务代码,边看文档边问最快。
一、分层架构

二、Risk 在 FP 中的位置
三、服务依赖

3.1 credit_risk_engine 内部结构
Task 按节点类型分,不按场景:
task_antifraud/task_feature/task_namelist/task_challenge/task_invoke_policy… 每类节点一份代码、所有场景共用;某场景的规则/特征不在 task 代码里,在 DAG 配置。正交结构:
task= 横向能力(workflow/task_*.go,场景无关)×场景= 纵向 DAG(Holmes 配,挑几种节点串起来)。一个场景就是选几个节点类型编成 DAG,每个节点复用对应 task。Task Registry(
flow_engine/flow_task_registry.go)=NodeType → TaskExecutable注册表,加全新节点类型才动它。
3.2 下游:怎么调 × 调谁
engine 是唯一的 hub,下游服务彼此不通信(唯一例外:nearline_engine → data_server 异步上报)。编排责任全部集中在 engine,是中心化设计。
| 调用形态 | 下游服务 | 职责 | 何时触发 |
|---|---|---|---|
| ① vendor 库 · 进程内 | 无对端 | 规则 DAG 执行 | 节点 anti_fraud@base |
| ② SDK 门面 → RPC | risk_feature_core | 特征计算 | 规则需要某个 feature |
| ② SDK 门面 → RPC | risk_list | 黑 / 白 / 灰名单 | 规则需要查名单 |
| ③ 直接 RPC | credit_risk_challenge | 下发 2FA 加验 | 反欺诈拿不定主意 |
| ③ 直接 RPC | risk_policy | 授信定价 | 需要授信判断 |
| ④ 数据层 · 不挂 task | credit_risk_data_server | 统一数据读写 | 贯穿全程,与节点无关 |
形态 ① 不产生下游服务 ——
risk-antifraud没有对端,纯计算跑在进程内(调用 watson/moriarty 规则引擎)
data_server不通过 task 触达 —— 它是 engine 的数据层,落库 / 查历史 / 结果上报时调
3.3 risk-antifraud 的规则引擎
Watson/Moriarty 只干一件事:
输入:一条规则表达式 + 一堆变量值
输出:命中 / 未命中
位置是三层嵌套,全在一个进程里 —— 它们是「依赖的依赖」,engine 不直接调:
credit_risk_engine(Go 进程)
└─ risk-antifraud(vendor 库)
└─ Watson / Moriarty(vendor 库)
Watson 的 parser 由 ANTLR 生成,文法是自研的。
Watson 和 Moriarty 不是分工,是新旧替换。靠 context 标志位切换。
四、DAG (Directed Acyclic Graph) 工作流
通过过将决策流程建模为 JSON 配置树(存储于 DB),运维和策略团队可通过 Holmes 平台在不发版的情况下修改执行逻辑(Graph)
admin 写 ----> MySQL ---> 分钟级刷新 ----> engine 内存缓存 ----> FlowEngine 读
策略同学在 antifraud-admin(Holmes)上编排 DAG,配置以 JSON 存进 mysql; credit_risk_engine 定期加载到内存(定时轮询刷新内存),运行时按 DAG 逐节点执行,每个节点对应一个 Task 实现,Task 内部按需调用下游服务
4.1 两层 DAG
外层 · 工作流 DAG
内层 · 决策流 DAG
以 anti_fraud 为例
两层配置结构一模一样(RootNode+Nodes+ 条件跳转),但跑在不同引擎里、配在不同的表里、各自独立灰度,可以单独放量。
监控对应:告警里的 sceneId 也分两层——工作流入口 sceneId = 外层;反欺诈节点 sceneId(
anti_fraud_<数字>)= 内层。用workflow_run_steps_stat的run_steps分辨,查错层查空。见 08-Risk 值班指南 与 02-Holmes 平台业务全景。
4.2 图的表示
没有独立的 edge 对象,边内嵌在节点里,
map[nodeId]Node, 靠 Goto 里的 id 引用串起来,遍历时现查现跳
"anti_fraud_10034": {
"NodeType": "anti_fraud@base",
"Goto": [ ← 出边列表,长在节点身上
{"Condition": "Output.CheckResult == 2", "NodeId": "return_reject"},
{"Condition": "ELSE", "NodeId": "challenge_10034"}
]
}图论上叫邻接表(adjacency list),不是边列表、也不是邻接矩阵。
| 维度 | 实现 |
|---|---|
| 节点定义 | Task Registry 注册 NodeType → TaskExecutable |
| 边定义 | 节点内的 Goto 列表 |
| 图表示 | map[nodeId]Node 邻接表 |
| 路由逻辑 | 表达式字符串(配置,非代码) |
| 状态 | Variable(map[string]interface{}) |
| 状态作用域 | 按节点分区:{nodeId}_Output |
| 编译 | 无编译期,运行时解释 JSON |
| 改流程 | 改 DB 配置 → 等缓存刷新 |
Trade-off
- 路由逻辑是数据不是代码 —— 策略同学改流程不用发版。
代价:Output.CheckReslt拼错编译期发现不了,run 起来才能发现 - 状态按节点分区 —— 分支条件能显式引用任意历史节点,不只是上一个:
"anti_fraud_10034_Output.CheckResult == 1 && condition_5_Output.Level > 2"
这是”图”而非”链”的实质体现。
4.3 驱动机制:一个 for 循环
DBFlowEngine.RunFlowInstance:
currentStep := configTree.RootNode // 重入时改用 CurrentStep
for {
node := configTree.Nodes[currentStep] // 邻接表查节点
task := taskRegistry.GetTask(node.NodeType) // NodeType → 实现
result, _ := task.Run(...)
allResult[currentStep] = result // 存这一步输出
runStep += currentStep + "@" + ver + "/" // 留痕
if AsyncFlag == 3344 { FINISHED = 0; return } // 挂起等重入
next := stepChange(node, result, allResult) // 算下一步
if next == "END" { FINISHED = 1; return }
currentStep = next
}单线程顺序执行,不并行。
4.4 stepChange —— 算下一个节点
func stepChange(curNode, preResult, allResult) string {
// 1. 构造求值环境
env := map[string]interface{}{}
for nodeId, v := range allResult {
env[nodeId+"_Output"] = v.ToMap() // 历史节点,按 id 分区
}
env["Output"] = preResult.ToMap() // 当前节点,固定叫 Output
// 2. 依次求值 Goto 分支,第一个 true 的生效
for _, branch := range curNode.Goto {
if ok, _ := expr.Eval(branch.Condition, env); ok {
return branch.NodeId
}
}
}求值环境(env) = 表达式里「变量名 → 值」的字典,等价于 Python eval(expr, env_dict) 的第二个参数。
每跑完一个节点 env 就多一项,越往后能引用的越多。
表达式引擎(antonmedv/expr):词法分析 → AST → 求值,AST 预编译缓存,不是每次重新 parse 字符串。
最后一条通常配 "Condition": "true" 兜底。没兜底又全不命中 → 路由失败。
4.5 工作流状态
Variable 底层是 map[string]interface{},整个流程共用的一个可变数据袋
| Key | 说明 |
|---|---|
INSTANCE_INPUT | 完整的 RiskGatewayReq(JSON) |
REQ_CONTEXT | 请求上下文(AppID / SceneID / ReqType) |
CURRENT_STEP / RUN_STEP | 当前节点 / 已执行路径 |
WORKFLOW_NO / WORKFLOW_ID | 流水号 FlowNo / 配置树 ID |
CHECK_POINT | 断点数据 |
{nodeId}_Output | 某节点的输出(JSON) |
{nodeId}_Var | 某节点的输出(Variable) |
WorkflowReentry | 是否重入 |
GLOBAL_FEATURES | 全局特征,跨节点传递 |
ProcessFeatureMap | 反欺诈 process 特征,跨 scene 用 |
ext_info / task_id | 当前节点的配置和 ID |
流程完成标志:FINISHED = 1 正常结束,FINISHED = 0 挂起中等重入。 | |
| 路径不是预先规划的,是一步步求值走出来的。 |
4.6 RunStep —— 执行轨迹
workflow@{workflowId}|node1@ver/node2@ver/.../END@ver
一眼知道请求实际走了哪条路、在哪停的。
4.7 节点类型(20+,摘主要)
| NodeType | 说明 |
|---|---|
anti_fraud@base | 反欺诈决策,可能返回异步 |
challenge_standard@base | 加验,含 OTP/LC/Popup/TapSecure 子流程 |
invoke_policy@base | 调授信,可能 pending |
condition@base | 条件判断(基于 process feature) |
anti_fraud_group / process / output@base | 反欺诈组 / 处理 / 输出 |
return_result@{version} | 结果返回,多版本按场景区分 |
otp / pin / tap_secure / sing_pass@base | 各类加验 |
liveness_check_init / facial_matching | 活体、人脸 |
注册键 = Name() + "@" + Version()。配置树里的 NodeType 查不到就跑不起来。 |
4.8 异步与重入
加验、授信 pending 等不了 → AsyncFlag == 3344 挂起,状态存 MySQL + Redis。
四种重入触发方式:
- 定时 Kafka 生产者
- Redis ZSET 延迟队列
- 授信结果回调(
InvokePolicyResultNotify) - 查询触发(
QueryAndReactivateAsyncRsp)
重入时从 workflow_tab / workflow_steps_tab 恢复 CheckPoint,Reentry=true,从 CurrentStep 接着跑。有分布式锁防并发重入。
4.9 配置树发布
路由表三档:GOLIVE / GRAY / SHADOW(正式 / 灰度 / 陪跑)。
新请求按路由表随机选 workflowId。配置存 MySQL,engine 内存缓存分钟级刷新 —— 改完不是立即生效。
4.10 工作流 DAG 示例
Credit-SPL-MP 支付:
graph LR START(["RootNode"]) --> AF["anti_fraud_10034<br/>anti_fraud@base"] AF -->|CheckResult == 2| REJ["return_reject"] AF -->|CheckResult == 1| IP["invoke_policy<br/>PolicySceneId: 20001"] AF -->|ELSE| CH["challenge_10034<br/>challenge_standard@base"] CH -->|result == 1| IP CH -->|result == 2| REJ IP -->|result == 1| PASS["return_pass"] IP -->|result == 2| REJ PASS --> END1(["END"]) REJ --> END2(["END"]) style AF fill:#ffd style CH fill:#fdd style IP fill:#dfd
五、请求生命周期(六阶段)
sequenceDiagram participant B as 上游业务 participant G as Gateway participant E as risk_engine participant A as risk_antifraud participant F as feature_core participant C as challenge participant P as risk_policy participant EA as EA网关/外部 B->>G: ① HTTP 请求 G->>E: 协议转换 → gRPC E->>E: ② 识别 appId+SceneId<br/>选定工作流 E->>A: ③ 执行规则 DAG<br/>(进程内调用,非 RPC) A->>F: 取特征 F->>EA: ⑥ 外部数据<br/>设备/IP/人脸/KYC EA-->>F: F-->>A: 特征值 A-->>E: pass / reject / 不确定 alt 不确定 E->>C: ④ 加验<br/>OTP/LC/FaceMatch/PIN C-->>E: 加验结果回调 end opt 需要授信 E->>P: ⑤ invoke_policy P-->>E: pass/reject/verify/pending end E-->>G: 最终决策 G-->>B: 响应 E-)+Kafka: 异步 → nearline_engine
六、输出形态
graph LR R["Risk 决策"] --> P["Pass<br/>继续业务"] R --> RJ["Reject<br/>拦截"] R --> C["Challenge / Verify<br/>补充验证"] R --> PD["Pending<br/>异步/人工推进"]
不是单一评分,是可执行的业务结论。
七、服务 ↔ 仓库映射
| 服务 | 仓库 | 职责 |
|---|---|---|
| credit_risk_engine | credit_risk_engine | 决策引擎,工作流调度,DAG 执行 |
| — (vendor 库) | risk-antifraud | 反欺诈引擎库,编译进 engine,非独立服务 |
| credit-dynamic-gateway | credit-unified-gateway | HTTP→gRPC,路由,限流熔断 |
| credit_risk_challenge | credit_risk_challenge | 加验:LC/OTP/PIN/TapSecure/FaceMatch |
| credit_risk_data_server | credit_risk_data_server | HBase/Redis/MySQL 读写,Kafka |
| risk_feature_core | risk_feature_core | 特征计算,gRPC 流式 |
| risk_feature_server | risk_feature | 在线特征查询 |
| risk_feature_task | risk_feature_task | 特征离线任务 |
| risk_feature_http_config_server | risk_feature_http_config_server | 特征 HTTP 配置服务 |
| risk_list | risk_list | 名单 CRUD 与匹配 |
| risk_nearline_engine | risk_nearline_engine | 消费 Kafka 结果,补全,回放 |
| antifraud_admin | antifraud-admin | Holmes 后台 |
| risk_aml | risk_aml | 反洗钱 |
其他相关仓库:
- credit_risk_gateway_proto —— PB 定义
- watson —— 旧版规则引擎
- risk-underwriting/moriarty —— 新版规则引擎(当前使用),在授信仓库下
- risk_af_agent_hub —— AI 能力底座