智能风控最难的地方,不是写出一条“同一设备十分钟登录五次”的规则,而是在数据乱序、消息重复、服务故障、规则频繁变化和业务高峰同时发生时,仍然做出稳定、可解释、可追溯的决策。
一个成熟系统需要同时回答四类问题:
- 现在是否放行:支付、登录、提现等请求通常要求毫秒级返回;
- 最近发生了什么:用户、设备、账户在滑动窗口内是否出现聚集或突变;
- 历史上是否异常:离线数据能否形成标签、基线、画像与训练样本;
- 为什么这样判断:每次决策能否还原使用的事件、特征、规则和模型版本。
因此,生产级风控不是一个 Spark 作业,也不是一个规则库。它是一套由同步决策、异步流计算、离线湖仓、策略平台和反馈闭环共同组成的决策系统。
总体架构:双路径、同事实、可回放
架构的关键是把“同步拦截”和“实时计算”分开:
- 同步在线路径:风险 API 读取在线特征,执行规则与模型,在严格延迟预算内返回通过、拒绝或人工审核;
- 异步实时路径:业务事件进入 Kafka,Spark Structured Streaming 以事件时间执行去重、滑动窗口、状态聚合和实时特征更新;
- 离线批处理路径:原始事实进入 HDFS 或对象存储上的湖仓,Spark Batch 负责回放、对账、特征回填、规则回测和模型样本构建;
- 统一策略与审计:规则、模型、特征定义和决策记录全部版本化,使线上结果可以离线复现。
这比让 Spark 直接阻塞每一个业务请求更稳妥。Structured Streaming 默认以微批方式执行,适合秒级状态计算;如果核心交易要求稳定的几十毫秒响应,同步决策服务必须独立部署,避免流处理调度、反压或 Checkpoint 抖动进入交易链路。
第一原则:建立不可变的风险事实流
Kafka 应被视为风险事件的事实日志,而不是临时消息管道。支付、登录、设备、账户变更、名单更新和人工审核结果都形成标准事件。
每条事件至少包含:
event_id 全局唯一,用于去重和追踪
event_type 事件类型及语义版本
event_time 业务实际发生时间
ingest_time 平台接收时间
entity_keys user_id / account_id / device_id / ip 等
payload 业务事实,不放计算后的临时判断
trace_id 贯穿请求、流任务与决策记录
schema_version 支持兼容演进
source 来源系统与环境
Topic 按业务事实和数据保留要求设计,不按每条规则建 Topic。分区键要匹配状态计算需要:账户规则使用 account_id 可以保持同一账户事件有序,但设备团伙规则又需要 device_id,因此常通过规范化事实流派生不同 keyed stream,而不是期待一个分区键解决所有问题。
生产者启用幂等和足够确认级别,Schema 使用 Avro、Protobuf 或带注册中心的 JSON Schema 管理兼容性。消费者仍然必须具备幂等能力,因为端到端是否“恰好一次”还取决于外部 Sink。风险系统更务实的目标是:至少一次传输 + 业务主键去重 + 幂等写入 + 可对账重放。
数据库变化不要依赖高频轮询。用户状态、账户等级、名单和商户信息可以通过 Debezium 等 CDC 进入 Kafka;业务数据库更新与领域事件必须一致时,使用 Transactional Outbox,避免“数据库成功、发消息失败”的双写裂缝。
用事件时间定义滑动窗口
风控窗口必须基于 event_time,而不是消息到达 Spark 的时间。网络延迟、移动端离线和 Kafka 积压都会造成乱序。如果用处理时间,重启或积压期间同一批业务事实可能得到不同结果。
典型窗口包括:
1 分钟窗口,每 10 秒滑动:支付瞬时爆发
10 分钟窗口,每 1 分钟滑动:账户失败次数与金额
24 小时窗口,每 15 分钟滑动:设备关联账户数
7 天状态:新旧收款方、常用地域和行为基线
窗口长度代表观察范围,滑动步长代表更新频率。步长越小,状态和计算开销越大。不要为了“更实时”把所有窗口设为每秒滑动;应从业务允许的发现延迟反推配置。
Watermark 定义系统愿意等待迟到数据多久,也决定状态何时可以清理。它不是“数据一定在这个时间内到达”的保证。应根据真实延迟分布制定,例如覆盖 99.9% 事件,再把超出水位线的数据写入迟到旁路,供补算、审计和指标监控,不能静默丢弃。
events = (
spark.readStream.format("kafka")
.option("subscribe", "risk.payment.v1")
.load()
.transform(parse_and_validate)
.withWatermark("event_time", "10 minutes")
.dropDuplicatesWithinWatermark(["event_id"])
)
velocity = (
events.groupBy(
window("event_time", "10 minutes", "1 minute"),
"account_id",
)
.agg(count("*").alias("tx_count_10m"),
sum("amount").alias("tx_amount_10m"))
)
示例表达的是语义,不应直接复制到生产。真实作业还需要 Schema 校验、坏数据隔离、状态大小控制、Checkpoint 独立目录、输入速率限制、监控和 Sink 幂等设计。
实时特征不能只存在 Spark 内存里
滑动窗口结果、最近设备集合、账户速度特征和风险计数需要写入低延迟在线存储,由同步决策服务读取。可以使用 Redis、Cassandra、HBase 或团队已有的高可用 KV,但必须先定义特征契约:
feature_name + entity_key + value
event_time + computed_at
definition_version + producer_job_version
ttl + freshness_sla
在线特征写入必须幂等。使用 foreachBatch 时,可将 query_id + batch_id 或业务窗口主键作为提交标识,采用 Upsert、事务表或 Commit Log 防止重试造成重复副作用。Spark 的 Checkpoint 能恢复源 Offset 和计算状态,但不会自动让任意外部数据库获得端到端 exactly-once。
还要定义陈旧特征策略。在线存储不可用或特征超过 freshness SLA 时,决策引擎应知道是“真实值为零”还是“数据缺失”,并根据风险等级选择保守规则、降级模型、人工审核或有限放行。
规则与模型不是竞争关系
成熟风控通常采用组合决策:
硬规则:监管、黑名单、明确禁止条件
速度规则:次数、金额、关联实体和时间窗口
模型分:欺诈概率、异常分、账户接管概率
策略编排:规则命中 + 模型分 + 业务成本 → 动作
规则适合确定性强、需要解释和立即响应的风险;模型适合多变量组合和难以手写的模式。最终动作不能只看一个分数,还要结合损失金额、客户价值、误杀成本与人工审核容量。
规则平台至少需要版本、状态、优先级、适用人群、生效时间、失效时间、作者、审批人和变更原因。发布流程采用:草稿 → 单元测试 → 历史回放 → Shadow → 小流量灰度 → 全量。规则命中必须记录规则版本与输入快照,但不要在高 QPS 主链路同步写重型审计数据库;先写可靠事件,再异步落库。
规则 DSL 应保持受限和可分析,避免允许任意脚本访问网络或数据库。规则数据和维表通过受控接口或预加载快照提供,不能让每条交易规则临时查询生产库。
离线层:从 Hadoop 文件堆升级为可治理湖仓
HDFS 仍适合本地 Hadoop 集群的大规模分布式存储;云上则通常使用对象存储。无论底层是什么,都不建议继续以裸 Parquet 目录作为唯一数据管理方式。使用 Iceberg 等开放表格式,可以获得原子提交、Schema 演进、分区演进、快照和 Time Travel,让回放与审计有稳定基础。
湖仓至少划分三层:
- 原始事实层:Kafka/CDC 原样归档,只追加,保留事件与 Schema 版本;
- 标准明细层:完成去重、主数据映射、隐私处理与统一时间语义;
- 特征与标签层:面向规则回测、模型训练、案件分析和指标统计。
Spark Batch 负责日终对账、历史窗口、特征回填、坏账与欺诈标签、策略效果评估。批处理与实时计算应共享特征定义或由同一套声明式逻辑生成,避免“训练时一个口径、线上另一个口径”。训练样本必须做 point-in-time join,只使用决策发生时已经可见的信息,防止未来数据泄漏。
流式写入 Iceberg 会产生大量 Snapshot 和小文件,需要独立维护任务定期压缩小文件、重写 Manifest、过期 Snapshot,并监控元数据增长。离线层不是无限保留:原始、特征、决策和敏感数据应分别设置保留期限与删除策略。
一次决策必须可以完整重放
每个决策至少留下以下 Decision Record:
decision_id / request_id / trace_id
event_id 与原始事实位置
特征值、特征时间和定义版本
命中规则及规则版本
模型名、模型版本、分数与阈值
最终动作、原因码和人工覆盖
决策时间、延迟、降级状态
审计不是只为合规,也用于回答生产问题:为什么同一客户昨天通过、今天拒绝?规则上线后误杀来自哪里?模型表现下降是数据漂移、特征延迟还是业务分布变化?
所有决策都应携带稳定的 reason code。面向运营的解释和面向客户的解释要分开设计,既保证可操作性,也避免暴露可被攻击者利用的规则细节。
高可用的关键是明确降级语义
系统不可能永不失败,生产设计的重点是失败时做什么:
| 故障 | 推荐策略 |
|---|---|
| Kafka 短时积压 | 在线路径继续使用带时间戳的最近特征,监控新鲜度并切换保守策略 |
| 在线特征库不可用 | 使用本地缓存与核心规则;高风险交易转人工或拒绝 |
| 模型服务超时 | 快速熔断,降级到规则和最近稳定模型,不串行重试拖垮请求 |
| 规则配置错误 | 立即回滚版本;高风险规则双人审批,发布前 Shadow |
| Structured Streaming 作业失败 | 从持久 Checkpoint 恢复;原始 Kafka 数据和湖仓归档支持回放 |
| 数据质量异常 | 隔离坏数据、冻结受影响特征、告警数据负责人,避免错误扩散 |
同步决策服务跨可用区部署,设置整体延迟预算和每个依赖的超时;Kafka、Checkpoint、在线存储和湖仓分别制定 RPO/RTO。不要把“重启成功”当作灾备,必须定期演练从 Offset、Checkpoint 和历史事实恢复。
可观测性要同时看系统、数据和决策
只监控 CPU 和延迟不足以运营风控。至少建立三组指标:
- 系统指标:QPS、p95/p99 延迟、错误率、Kafka Lag、批次耗时、状态大小、Checkpoint 失败;
- 数据指标:事件量、Schema 失败、重复率、迟到率、空值率、特征新鲜度、实时离线一致性;
- 风险指标:通过/拒绝/审核率、规则命中率、模型分布、误杀率、欺诈损失、审核积压和策略收益。
告警要关联影响而非只关联阈值。例如 Kafka Lag 上升并不可怕,真正重要的是哪些特征已经超过 SLA、哪些决策正在降级,以及受影响的交易量。
从 MVP 到生产的落地顺序
阶段一:建立事实与最小决策闭环
- 统一风险事件 Schema、唯一 ID、事件时间和 Reason Code;
- Kafka 接入一种核心事件,原始数据同步归档湖仓;
- 上线独立决策 API、少量硬规则和完整 Decision Record;
- 建立延迟、错误率、Lag 和决策分布基线。
阶段二:引入实时状态
- 用 Structured Streaming 实现 2–3 个高价值滑动窗口;
- 建立 Watermark、迟到旁路、去重和幂等 Sink;
- 接入在线特征存储并定义 freshness SLA;
- 演练作业重启、Kafka 积压和特征库降级。
阶段三:规则平台与离线回放
- 规则版本化、审批、Shadow、灰度和一键回滚;
- 建立 Iceberg 明细、特征、标签和决策表;
- 对新规则做历史回测并校验实时/离线特征一致性;
- 引入人工审核反馈和案件闭环。
阶段四:模型化与持续优化
- 通过 point-in-time join 建立无泄漏训练集;
- 模型注册、可解释输出、Shadow、灰度和漂移监控;
- 用损失、误杀和运营成本共同优化阈值;
- 定期复盘规则重叠、无效特征、技术债和灾备能力。
最后的判断标准
生产级智能风控不等于“实时技术足够多”。真正成熟的系统应该具备:
- 同步决策延迟可控,流计算不会拖垮交易主链路;
- Kafka 保存可重放事实,重复、乱序和迟到都有明确语义;
- 实时与离线共享特征口径,规则和模型全部版本化;
- 每次决策可以解释、审计和复现;
- 任一依赖失败时,系统知道如何降级,而不是随机失败;
- 策略效果最终由欺诈损失、误杀成本和客户体验衡量。
风控的本质不是尽可能拒绝风险,而是在不断变化的信息和有限时间内,做出成本可控、证据充分、责任清晰的决策。Kafka、Spark、Hadoop 和规则引擎只是基础设施;真正的核心,是把实时事实、历史知识和业务判断组织成一条持续学习的决策闭环。