Administrator
发布于 2025-04-16 / 1622 阅读
40

流式处理在实时 AI 场景的应用

风控要在 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 有三个消费者,各干各的:

  1. 风控决策服务:消费后更新本地缓存(Caffeine),决策时优先查本地,查不到再查 Redis;
  2. 样本回流作业:把特征 + 推理结果 + 后续的人工审核结论拼成训练样本,落到样本库。这个很重要,模型上线两个月后我们就是靠这批样本做了第一次迭代;
  3. 监控作业:统计分数分布,发现分布突变说明特征或模型出了问题。

第 2 点值得单独说:AI 系统的数据飞轮必须在一开始就设计好。我们第一版没做样本回流,等想迭代模型时发现没有标注数据,只能回头补,多花了两周。现在样本回流是默认流程,每天回流量 180 万条。

效果数据

上线三个月后的对比(规则引擎 → 规则 + 模型):

指标纯规则规则 + 模型
欺诈订单召回率61.3%87.9%
误拦截率1.82%0.71%
决策 P9935ms167ms
特征计算延迟(端到端)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,那套集群的运维成本不低。

参考