# 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 内部结构

credit_risk_engine gitlab 链接

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 门面 → RPCrisk_feature_core特征计算规则需要某个 feature
② SDK 门面 → RPCrisk_list黑 / 白 / 灰名单规则需要查名单
③ 直接 RPCcredit_risk_challenge下发 2FA 加验反欺诈拿不定主意
③ 直接 RPCrisk_policy授信定价需要授信判断
④ 数据层 · 不挂 taskcredit_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 = 外层;反欺诈节点 sceneIdanti_fraud_<数字>)= 内层。用 workflow_run_steps_statrun_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 邻接表
路由逻辑表达式字符串(配置,非代码)
状态Variablemap[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。

四种重入触发方式:

  1. 定时 Kafka 生产者
  2. Redis ZSET 延迟队列
  3. 授信结果回调(InvokePolicyResultNotify
  4. 查询触发(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_enginecredit_risk_engine决策引擎,工作流调度,DAG 执行
— (vendor 库)risk-antifraud反欺诈引擎库,编译进 engine,非独立服务
credit-dynamic-gatewaycredit-unified-gatewayHTTP→gRPC,路由,限流熔断
credit_risk_challengecredit_risk_challenge加验:LC/OTP/PIN/TapSecure/FaceMatch
credit_risk_data_servercredit_risk_data_serverHBase/Redis/MySQL 读写,Kafka
risk_feature_corerisk_feature_core特征计算,gRPC 流式
risk_feature_serverrisk_feature在线特征查询
risk_feature_taskrisk_feature_task特征离线任务
risk_feature_http_config_serverrisk_feature_http_config_server特征 HTTP 配置服务
risk_listrisk_list名单 CRUD 与匹配
risk_nearline_enginerisk_nearline_engine消费 Kafka 结果,补全,回放
antifraud_adminantifraud-adminHolmes 后台
risk_amlrisk_aml反洗钱

其他相关仓库: