风控要在 200ms 内给出决策,而模型要 150ms
我们电商风控原来是一套 Drools 规则引擎,几百条规则,命中就拦截。去年底开始加模型:用用户最近 5 分钟的行为序列算实时特征,喂给一个轻量模型打分,超过阈值就拦截。
难在时间预算。风控决策必须在用户下单后 200ms 内返回,而模型推理本身要 120~180ms。留给特征计算的时间不到 40ms,传统的"下单时去查各种表聚合"根本来不及。所以特征必须提前算好,下单时只做读取。
这就是流式处理进场的地方。这篇记录我们用 Flink + Kafka 做实时特征计算和推理结果回写的过程。
整体链路
最终跑起来的链路是这样:
用户行为埋点 ──┐
├──> Kafka(behavior) ──> Flink 实时特征作业 ──> Redis(特征库)
订单事件 ──────┘ │
│ 触发
▼
Kafka(infer.req) ──> 推理服务
│
▼
Kafka(infer.result)
│
┌─────────────────────────────┼──────────────┐
▼ ▼ ▼
风控决策(查 Redis 取结果) Flink 回写作业 监控/样本库
这里最关键的设计是推理是异步的。下单事件进 Kafka 后,Flink 作业触发一次异步推理,结果写回另一个 topic 并落 Redis(带 TTL)。风控决策服务在需要做决策时,直接查 Redis 拿这个订单的推理结果——如果还没算出来(超时未就绪),就走规则兜底。
为什么不让决策服务同步调模型?两个原因:一是 120~180ms 的推理时间会把决策服务的 P99 直接拉满;二是同一个用户短时间内可能触发多次风控检查(浏览、加购、下单),异步化后可以合并掉重复请求,我们实测合并了 34% 的冗余推理。
为什么选 Flink 而不是 Kafka Streams
我们两个都试过。Kafka Streams 的好处是轻,就是个 Java 库,不用单独部署集群,我们一开始用的它。但后来遇到两个问题换成了 Flink:
- 状态太大:我们要维护每个用户 24 小时的滑动窗口行为,峰值 800 万活跃用户,状态总量到了 40GB 左右。Kafka Streams 的 RocksDB 状态存在本地磁盘,扩缩容时状态迁移很慢(实测一次扩容迁移花了 23 分钟);
- 窗口语义不够用:我们需要会话窗口和自定义触发器,Kafka Streams 虽然能做但写起来很绕。
Flink 的 checkpoint 机制和成熟的 savepoint 迁移在这两点上明显更好用。代价是要维护一个 Flink 集群(我们用的 K8s 上跑的 Flink 1.19,JobManager 2 核、TaskManager 16 个 × 4 核)。
如果你们的状态量在 10GB 以内、窗口逻辑简单,Kafka Streams 其实是更省心的选择。别一上来就上 Flink。
实时特征计算
核心作业长这样,用的是 Flink DataStream API + KeyedProcessFunction:
DataStream<BehaviorEvent> events = env
.fromSource(kafkaSource, WatermarkStrategy
.<BehaviorEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((e, ts) -> e.eventTime()),
"behavior-source");
DataStream<UserFeatures> features = events
.keyBy(BehaviorEvent::userId)
.process(new FeatureAggregator())
.name("feature-agg")
.uid("feature-agg");
features.sinkTo(redisSink).name("feature-sink");
特征聚合函数里用 MapState 存用户的时序行为:
public class FeatureAggregator
extends KeyedProcessFunction<String, BehaviorEvent, UserFeatures> {
private transient MapState<Long, BehaviorEvent> window;
@Override
public void open(Configuration cfg) {
MapStateDescriptor<Long, BehaviorEvent> desc =
new MapStateDescriptor<>("win", Types.LONG,
TypeInformation.of(BehaviorEvent.class));
window = getRuntimeContext().getMapState(desc);
}
@Override
public void processElement(BehaviorEvent e, Context ctx,
Collector<UserFeatures> out) throws Exception {
window.put(e.eventTime(), e);
// 注册 5 分钟后的清理定时器,防止状态无限膨胀
ctx.timerService().registerEventTimeTimer(
e.eventTime() + Duration.ofMinutes(5).toMillis());
out.collect(compute(e.userId()));
}
@Override
public void onTimer(long ts, OnTimerContext ctx,
Collector<UserFeatures> out) throws Exception {
long cutoff = ts - Duration.ofMinutes(5).toMillis();
window.iterator().forEachRemaining(entry -> {
if (entry.getKey() < cutoff) window.remove(entry.getKey());
});
}
}
这里有个必须注意的点:一定要注册清理定时器。第一版我们没写 onTimer,跑了一周状态从 6GB 涨到 40GB,checkpoint 从 12 秒变成 4 分多钟,最后作业因为 checkpoint 超时反复重启。加上定时器清理后状态稳定在 7~9GB。
我们算的特征大概 40 个,主要是:
- 时间窗口类:5 分钟内点击次数、加购次数、搜索次数、切换收货地址次数;
- 会话类:本次会话时长、浏览商品类目数、价格区间跨度;
- 对比类:当前收货地址与历史常用地址的距离、下单金额与该用户历史均值的偏离度。
乱序数据的处理
设了 5 秒的乱序容忍,但总有些数据迟到更多——移动端弱网环境下,我们观察到最长有 40 多秒的延迟。迟到数据直接丢弃会影响特征准确性,我们的处理是用侧输出流收集迟到的事件,单独处理:
OutputTag<BehaviorEvent> lateTag =
new OutputTag<BehaviorEvent>("late-events") {};
SingleOutputStreamOperator<UserFeatures> main = events
.keyBy(BehaviorEvent::userId)
.process(new FeatureAggregator())
.sideOutputLateData(lateTag); // 需要 window 算子
// 迟到数据:不重算,只更新一个"迟到修正"标记,供下游参考
main.getSideOutput(lateTag)
.addSink(new LateEventCorrector());
实测迟到率约 0.3%,影响有限。我们没有为这 0.3% 引入复杂的重算逻辑,只是在特征里加了个 dataCompleteness 字段,迟到严重的时候模型会收到这个信号。
推理结果回写
推理服务是独立的 Spring Boot 应用,消费 infer.req,调模型,结果写 infer.result。看起来简单,实际有几个坑。
异步推理的幂等和去重
同一个订单可能因为上游重试被多次投递。我们按 orderId 做幂等,Redis 里存"推理中"标记,TTL 10 秒:
@KafkaListener(topics = "infer.req", concurrency = "12")
public void handle(InferRequest req) {
String lockKey = "infer:lock:" + req.orderId();
Boolean acquired = redis.opsForValue()
.setIfAbsent(lockKey, "1", Duration.ofSeconds(10));
if (!Boolean.TRUE.equals(acquired)) {
return; // 同一订单的推理正在进行,丢弃重复请求
}
try {
RiskScore score = model.predict(req.features());
redis.opsForValue().set("infer:result:" + req.orderId(),
toJson(score), Duration.ofMinutes(30));
kafkaTemplate.send("infer.result",
req.orderId(), new InferResult(req.orderId(), score));
} finally {
redis.delete(lockKey);
}
}
这里 concurrency = "12" 配合手动 ack,保证并发消费。注意幂等锁的 TTL 必须大于最长推理时间——我们的模型 P99 是 180ms,设 10 秒绰绰有余,但如果模型卡住,10 秒后锁自动释放会允许重复推理,这个我们接受了(重复推理只是浪费钱,不影响正确性)。
背压:模型扛不住的时候
这是上线后遇到的最大问题。推理服务的消费速度跟不上生产速度时,Kafka lag 会持续上涨。我们一开始没做限制,消费者拼命拉消息,结果大量请求同时在等 GPU,反而让每个请求的延迟都变长,形成恶性循环。
解决办法是在消费端限制并发 + 主动降速:
semaphore.acquire(); // 限并发,等于 GPU 能承受的最大并发
try {
RiskScore score = model.predict(req.features());
} finally {
semaphore.release();
}
同时把 max.poll.records 从默认 500 调到 20,配合手动提交 offset。lag 上涨时监控告警,我们人工扩容。这里没有做自动扩缩容,因为 GPU 机器扩容需要几分钟,自动扩缩容的收益不明显,反而增加了复杂度。
结果回写触发的下游动作
infer.result 这个 topic 有三个消费者,各干各的:
- 风控决策服务:消费后更新本地缓存(Caffeine),决策时优先查本地,查不到再查 Redis;
- 样本回流作业:把特征 + 推理结果 + 后续的人工审核结论拼成训练样本,落到样本库。这个很重要,模型上线两个月后我们就是靠这批样本做了第一次迭代;
- 监控作业:统计分数分布,发现分布突变说明特征或模型出了问题。
第 2 点值得单独说:AI 系统的数据飞轮必须在一开始就设计好。我们第一版没做样本回流,等想迭代模型时发现没有标注数据,只能回头补,多花了两周。现在样本回流是默认流程,每天回流量 180 万条。
效果数据
上线三个月后的对比(规则引擎 → 规则 + 模型):
| 指标 | 纯规则 | 规则 + 模型 |
|---|---|---|
| 欺诈订单召回率 | 61.3% | 87.9% |
| 误拦截率 | 1.82% | 0.71% |
| 决策 P99 | 35ms | 167ms |
| 特征计算延迟(端到端) | — | P99 890ms |
| 推理服务 QPS | — | 峰值 4200 |
| 月均资损 | ¥87 万 | ¥23 万 |
决策 P99 从 35ms 涨到 167ms,在 200ms 预算内,但余量不大。这也是为什么决策服务查本地缓存而不是每次查 Redis——本地缓存命中率 92%,命中时延迟只有 8ms。
踩过的坑清单
- 状态没清理:上面说过了,checkpoint 从 12s 涨到 4 分钟;
- checkpoint 间隔设太短:一开始设了 10 秒,40GB 状态下每次 checkpoint 都要几分钟,导致作业一直在做 checkpoint。改成 3 分钟后正常;
- RocksDB 本地磁盘打满:Flink 的 RocksDB 状态存在 TaskManager 本地盘,我们给的 100GB 云盘被写满过一次。现在监控里加了磁盘水位告警,水位 70% 就处理;
- Redis 热 key:某个大促期间,一个头部网红店铺的用户行为特征被高频访问,形成热 key。后来对特征按 userId 做了分片,并且加了本地缓存;
- 特征穿越:这个最隐蔽。我们在拼接训练样本时,不小心把"订单已完成"这个事后才能知道的状态算进了特征,模型离线 AUC 高达 0.97,上线后只有 0.78。排查了一周才发现。离线指标异常好的时候,先怀疑特征穿越。
小结
实时 AI 场景里,流式处理的角色是"把模型需要的数据提前准备好",以及"把模型产出的结果送到该去的地方"。它本身不做智能的事,但没有它,实时 AI 根本跑不起来。
几个我觉得最重要的经验:状态一定要设计清理机制,否则迟早爆炸;推理必须异步化,同步调用会拖垮主链路;样本回流要在第一版就做,不然后面补代价很大;离线指标好得离谱时,先怀疑特征穿越而不是庆祝。
最后是选型建议:状态量小、逻辑简单用 Kafka Streams,省运维;状态量大、需要精确语义和灵活窗口用 Flink。别为了"技术先进性"上 Flink,那套集群的运维成本不低。