
做了这么久实时数据处理我一直有个很深的体会AI原生应用能不能真正聪明很多时候不取决于模型有多先进而是取决于数据到决策的时间差有多短。所谓实时决策简单说就是让系统在海量数据流动的过程中快速识别出值得采取行动的模式。这个场景在风控反欺诈、个性化推荐、实时竞价、智能调度、工业预测性维护里面都特别常见。比如用户刚点了一下商品系统要在几百毫秒内判断要不要给他推荐同类产品一笔交易刚发生模型要马上判断这单是不是有风险。这些需求催生出一个核心的技术底座——流处理技术。我打算在这篇里把AI原生应用如何利用流处理做实时决策这件事从原理、选型到实战坑位完整捋一遍。这篇文章适合谁看适合正在做推荐系统、风控系统、实时数据平台或者准备把AI能力从离线batch搬到在线场景的工程师。哪怕你之前没用过流处理框架只要能看懂基本的编程概念也能顺着这篇把整个技术脉络摸清楚。1. 实时决策碰到的第一个硬墙延迟的归因很多团队一开始做实时决策都会先遭遇一次认知颠覆。你以为模型的响应速度是瓶颈结果调了半天发现数据链路才是真正的拖油瓶。1.1 传统离线路径的短板传统的架构大多是这样的业务数据落库然后离线任务每天凌晨跑一次或者每小时跑一次把特征算好存进特征库模型推理也是针对一批数据统一处理最后把结果写回线上。在这个模式下从一条数据产生到模型决策生效延迟通常是以小时甚至天来计算的。这套模式在用户量稳定、业务规则简单的年代没有什么大问题。但当竞争激烈起来业务方提的需求就变了。风控这边说我们用昨天的数据训练出来的模型识别今天的新型欺诈基本等于用去年的地图找今年的路推荐那边说用户现在正在看的这个视频如果我们能知道他刚刚点赞了什么转化率就能再上一个台阶运营那边更直接用户已经把你App卸载了你第二天才推送挽回优惠券有意义吗说白了这些需求的核心痛点只有一个决策响应速度跟不上数据产生的速度。1.2 流处理和快速决策之间的真实关系流处理要解决的不是把数据变多而是把数据产生到数据可用之间的时间窗口压缩到秒级甚至毫秒级。它的核心思想是数据一旦产生就立即进入处理管线算完结果马上可以供决策系统调用。我见过不少团队在这里有个误解以为只要上了Kafka把数据实时地搬到消息队列里就算实时了。实际上Kafka只是传输通道真正让数据活起来的是传输之后的那一连串处理逻辑——数据清洗、特征拼接、窗口聚合、模型打分、结果下发。这一整套东西才是流处理的本来面貌。为什么说AI原生应用尤其依赖这套能力因为AI模型的决策质量极度依赖特征的时效性。一个用了过期特征的模型哪怕模型结构再精巧也只是一个精致的算盘。流处理技术恰恰能把特征新鲜度这个问题从根本上解决掉。所以从一开始实时决策的技术选型就没有绕路走直接奔着流处理这个方向来了。2. 选型之前先看懂框架的脾气四个主流流引擎的取舍技术选型是个很磨人的环节因为每个框架都有它的性格。没有绝对的好坏只有跟你的场景匹配不匹配。我梳理一下市面上几款主流流处理引擎各自适合什么活、干不了的活是什么都摊开说清楚。2.1 候选框架画像对比首先是Apache Flink。它是我见过在有状态的流处理这件事上做得最彻底的一个引擎。它提供真正的逐条事件处理而不是微批天然就适合需要低延迟、复杂事件处理、精确一次语义的场景。Flink的状态管理能力尤其突出可以在内存或者RocksDB里维护跨事件的聚合状态配合Checkpoint机制实现故障恢复。社区活跃度极高基本已经成了实时计算领域的事实标准。然后是Apache Spark Structured Streaming。这一套的底层是微批处理把连续的数据流切成一个个小批次来处理。它最大的优点是和Spark生态无缝衔接代码写起来舒服适合从离线批处理迁移过来的团队延迟能做到秒级。但它在真正的毫秒级场景下有些吃力如果业务要求极低延迟会感觉有点使不上劲。再聊聊Kafka Streams。它不是一个独立的集群而是一个Java库直接嵌在你的应用进程里。好处是部署简单和Kafka的集成是原生的适合那些不想引入一套独立计算集群的团队。但它的状态管理和窗口能力相对弱一些更适合简单的流式ETL和事件路由。还有一个后起之秀RisingWave它主打的是云原生流数据库直接用SQL来做流处理底层的存储和计算是分离的。它的特点是上手门槛极低——如果你会写SQL基本上十分钟就能开始写流任务。适合数据团队以SQL为主要开发语言、想快速搭建实时数仓或实时特征平台的场景。我之前在一个项目里就是用它做实时特征拼接比用Flink写Java代码省了大概一半的工程量。引擎处理模式延迟量级状态管理适合场景Flink逐条事件毫秒强RocksDB/内存复杂事件处理、风控、精准一次Spark Streaming微批秒级一般从离线迁移、ETL为主Kafka Streams逐条事件毫秒中等轻量流式应用、Kafka生态内闭环RisingWaveSQL流式秒级内云原生存储实时特征库、实时数仓快速构建2.2 我们最终的选择和依据我参与的一个实时风控项目最终选了Flink作为主力计算引擎。原因其实很简单风控场景里一次事件可能要在秒级窗口内关联几十个历史事件还要保持精确一次的处理语义——如果数据处理重复或丢失决策结果的可靠性就无从谈起。Flink的状态后端、事件时间处理、分布式快照这套组合拳给了我足够的信心。同时我把一部分实时特征拼接的活交给了RisingWave因为它对SQL的支持太友好了数仓团队的人几乎不需要重新学一套API。这里面有个很重要的思路实时决策架构不一定要押注在一个引擎上可以按场景拆开让不同的引擎干自己最擅长的那件事。避免用一把锤子敲所有的钉子。2.3 选型之外部署形态与运维成本选型的时候大家容易忽略部署形态带来的隐性成本。Flink如果你用Standalone模式光高可用就要自己处理用Flink on Kubernetes倒是省心但那套Operator配置起来又是一堆活。Spark Streaming如果集群里本来就有Spark那几乎零额外成本。Kafka Streams不需要独立集群运维最轻但如果你要把状态做大进程故障恢复和状态调度的复杂度会转嫁到自己身上。我的建议是如果团队规模小、没有专职实时计算运维优先考虑RisingWave或Kafka Streams这类简化运维的方案如果团队本就有大数据平台能力那Flink基本是必选项。3. 把流处理管线拆开看数据接入、加工与状态管理选完引擎接下来就是真刀真枪地搭管线了。一条标准的实时决策数据管线从上游到下游大致可以分为接入层、加工层、状态管理层和输出层。每一层都有几个容易被忽视的细节。3.1 接入层的三个致命细节第一消息队列的Topic分区数要和服务实例数对齐。很多新手把Kafka Topic建了12个分区结果Flink并行度只设了3一部分分区永远闲着数据全部堆在一个Subtask上延迟直接拉高。反过来如果并行度过大下游数据库写入压力会突然飙升。这个对齐的过程需要根据实际数据吞吐量反复调整。第二序列化格式的选择比想象中更重要。我们用过JSON和Avro做过对比同样的数据量JSON的解析耗时大概是Avro的三倍而且体积大两倍左右。对于实时性要求高的场景建议优先上Avro或者Protobuf配合Schema Registry管理版本这样上游字段变更时下游不会被直接冲垮。第三接入层一定要做流量整形。真实世界的流量从来没平缓过大促、热点事件、恶意攻击任何一阵涌流都可能把下游打垮。在接入层加一层限流和削峰填谷比如用Kafka自身的背压机制配合适当调整消费者的拉取速率能保护整条管线不至于因为峰值流量直接雪崩。3.2 窗口状态的实际操作流处理里最常用到的一个能力是维护一段时间内的累计状态。举个例子要判断某个用户在一分钟内点击了多少次商品如果没有窗口你要反复跨事件地去查数据库那效率太低了。窗口机制让你可以在内存中维护这个计数每个事件来的时候只要更新计数器就行。Flink提供了滚动窗口Tumbling Window、滑动窗口Sliding Window和会话窗口Session Window三种基本形态。我拿电商场景举个实际的例子滚动窗口每整分钟统计一次订单金额。适合做周期性的全量指标计算。滑动窗口每5秒输出一个持续了1分钟的滑动统计。适合做最近一分钟的行为轨迹这种需要连续刷新的指标。会话窗口用户连续30秒无动作才算会话结束。适合用户一次逛店行为的分析。状态管理的过程中最核心的配置是状态后端和TTL。状态后端建议直接选RocksDB除非你的状态真的非常小因为堆内存状态一旦数据量大起来GC问题会逼得你想骂人。TTLTime To Live一定不要省给Keyed State设置合理的过期时间比如状态超过24小时没更新就自动清掉否则累积的无用状态会拖垮整个作业。3.3 实时特征工程怎么做流处理里的特征工程和离线差不多但有三个关键差异一是特征只能看到当前时刻之前的数据不能look-forward二是特征的时效性要求极高比如最近5分钟内次数这种特征窗口一过期就作废三是特征需要支持在线拼接。我常用的做法是把特征分为两类一类是长周期聚合特征如用户近7天购买总额这类可以用离线Batch任务每天算一次写进特征存储另一类是短周期实时特征如用户最近5秒的滑动行为计数这类放在流处理里实时算。在线推理的时候把两类特征拼接起来一起喂给模型这样既不牺牲实时性又控制了计算成本。4. 时间语义决定了决策的敢不敢信流处理里最容易把人绕晕的是时间语义。我见过不止一个团队在这个问题上栽了跟头最后出来的结果数据看起来有模有样但实际上全都是错的。这一节我会把这个概念彻底讲明白。4.1 三套时间体系你要分清流处理框架里面有三个时间概念事件时间Event Time事件实际发生的时间由业务系统在数据里打上时间戳。这才是决策真正应该依赖的时间。摄入时间Ingestion Time数据进入流处理引擎的时间是引擎收到消息的时机。处理时间Processing Time数据在您的算子中真正被计算的那一刻的时间。刚开始做实时的时候我图省事直接用处理时间因为代码写起来最简单。后来发现问题很严重处理时间受系统负载影响队列一堆积处理时间就整体往后偏移了。同一个业务事件上午九点发生凌晨两点处理窗口统计的归属完全不同这要拿去喂模型等于把所有特征都对齐到错误的时区。所以实时决策场景下一律使用事件时间这一点没有任何妥协的余地。业务端必须在上游数据里写入可靠的事件时间戳最好精确到毫秒。4.2 Watermark机制在乱序世界里等待迟到数据现实世界中数据到达的顺序几乎不可能和事件发生的顺序完全一致。网络抖动、上游重试、应用异步写入都会导致先发生后到的情况。水位线Watermark就是用来应对这种乱序的。Watermark的含义可以简单理解成一个承诺它表示事件时间早于这个时间戳的所有数据我都已经接收完毕了。比如Watermark max_event_time - 30s意味着在计算窗口时系统只会触发那些事件时间低于水位线的窗口给晚到的数据留出最多30秒的等待时间。这里有个取舍问题等待时间设得太长窗口结果的产出就延迟了实时性下降等待时间太短迟到的数据会被丢弃准确性下降。我们在一个实时推荐的项目里把Watermark容忍度设成了10秒既满足了业务结果秒级可用的需求又把迟到的数据比例控制在了一个可接受的范围内。这个值没有银弹必须根据你的上游数据质量实测得出。需要注意的一点是迟到的数据是可以通过侧输出流Side Output接住的你可以把这些迟到数据导到另一个写进存储的Sink里用于事后分析。这样就兼顾了实时决策的低延迟需求和数据完整性的保障。4.3 窗口触发时机的细节Flink里窗口的触发条件是这样的当Watermark越过窗口的结束时间时这个窗口的结果会被正式计算并下发。所以你在选事件时间的同时还要去精确设计Watermark的生成频率和容忍度。举个例子假设一个1分钟的滚动窗口窗口范围是10:00:00到10:00:59。如果Watermark现在推进到了10:01:00那么这个窗口就会触发。如果一条事件时间10:00:55的数据在10:01:10才到它仍然会被Watermark放行进入这个窗口因为水位线已经越过窗口末端了。但如果这条数据是10:00:50而Watermark已经推进到10:01:20那它就会被判定为过迟数据不再进入窗口。理解这个细节你才能解释为什么有些统计结果和你在数据库里手动查出来的对不上。5. 让模型参与实时决策推理闭环的几种姿势流处理把数据算好了接下来面临一个关键问题怎么让AI模型在这条管道上发挥价值这里有很多工程师把流处理和模型推理当成两个孤立的系统来做结果就是数据和模型之间永远隔了一层。更好的方案是让模型直接嵌入到流处理管线的算子中。5.1 特征拼接与模型服务之间的小心机模型推理一般需要一个特征向量。这些特征一部分来自流处理算出来的实时特征一部分来自存储在离线特征库里的长期特征。实时决策系统在调用模型之前需要把这两部分拼接成一个完整的特征向量。这里的实践要点是不要把实时特征和离线特征放在两个地方同时拉取否则会引入一次额外的网络IO延迟。更常见的做法是把实时特征先写到共享的实时特征存储比如Redis或者RisingWave的物化视图模型服务统一从一个地方读取完整特征向量。这样既避免了分布式拼接的麻烦也方便多个模型复用同一套特征。模型推理本身可以通过RPC的方式调用独立的模型服务这个方案的优点是模型可以独立扩展但缺点是每次推理开销多一跳网络。如果单次推理本身很快比如几毫秒那这一跳网络的时间占比会很明显。更极致的方案是把轻量模型直接以函数的形式挂在流处理的UDF里模型跑在业务数据所在的地方避免数据传输的额外延迟。我在一个广告点击率预估的项目里就用了Flink的UDF直接加载一个几百KB的树模型实测单次推理耗时从8毫秒降到了2毫秒左右效果非常明显。5.2 实时反馈与模型更新循环实时决策系统的闭环还包含一个容易被忽略的部分决策结果的反馈回路。模型做出一个预测后这个预测是否准确最终会随着用户的反馈比如点击了没有、交易有没有欺诈在后续数据中体现出来。这些反馈数据要以实时或者近实时的方式回流到训练管线里面。最常用的做法是把决策结果和实际反馈组成新的样本写到一个实时训练样本通道中每隔一段时间触发一次增量训练或在线学习。这样模型才能保持对最新模式的敏感度。流处理在这里的价值是它同时兼任了在线推理和训练样本采集两件活用一个管道把两者统一起来。5.3 兜底和降级策略最后必须说的是任何实时系统都会故障模型推理也不是永远成功的。我们在生产环境里给决策服务配置了一套降级策略当模型服务响应超时或不可用时系统自动切换到规则引擎用预设的硬规则兜底。比如风控场景里如果模型服务挂了那么单笔金额超过阈值的交易一律转人工审核推荐场景里模型挂了就推荐热榜内容。这套降级策略的实现恰恰也要靠流处理——在算子层面捕获模型调用异常然后路由到兜底逻辑。没有这层保护一次模型服务抖动就可能导致整个决策管道堵塞那带来的损失可比用笨办法兜底大多了。6. 实战中绕不开的那些坑一条排查链路的完整复盘最后这部分我想用一次真实的排障经历把流处理落地过程中那些看着不起眼、却能让你翻车半天的坑讲透。那次我们搭好了一条实时风控管线跑了一周挺顺畅的然后就在某天下午突然发现决策产出延迟骤增。6.1 第一步从监控指标圈定问题面我先看的不是日志而是三个最核心的监控指标Kafka消费者组的Lag堆积消息数、Flink作业各算子的处理延迟Processing Delay、以及窗口触发的Watermark推进速度。三个指标一对照问题面很快就浮出来了Kafka Lag在持续上升说明消费者处理速度已经跟不上生产速度Flink的处理延迟指标也在同步拉高但是Watermark推进速度没有异常。这说明数据还在源源不断地进来卡点是出在计算本身而不是上游数据源或者网络传输。6.2 第二步揪出背压的真实源头接下来的排查重点是反压Backpressure。打开Flink Web UI一眼看到某个关键算子正好处于反压状态——它的输入还是通畅的但输出很慢说明问题的源头在下游算子。点开下游算子的详情一看这个算子正好是负责做业务规则判断的里面有一个很不起眼的操作查外部Redis。问题几乎瞬间就定位了。我们的规则判断里每来一条交易事件都要去Redis里查一下这个商户的历史风险等级。数据量小的时候这个查询平均耗时才1毫秒左右感知不到。但那天下午业务做了一次营销活动每秒进入的事件量突然暴涨到平时的十倍每个事件都要同步等待Redis查询的结果试图直接把外部查询吊死在算子链路里Redis的压力也因此被放大。最终整条管线的吞吐量就垮了。这个案例是个典型的流处理反模式在流计算算子内部同步调用外部存储。这会让外部系统的延迟问题和脆弱性直接传导到流处理管线上。正确的做法是什么把外部数据预加载进来作为维表关联比如Flink的Lookup Join配合缓存或者干脆把热点数据以广播状态的形式分发给所有算子实例让每个实例本地持有这份维表查询时完全不走网络。经过改造后同流量下平均处理时间从原来的几十毫秒降到了几毫秒Kafka的Lag很快就被消掉了。6.3 第三步顺带解决状态膨胀的隐患排查的过程中我还发现另一个隐患作业里注册的Keyed StateKey是用户IDTTL设了两周。这在正常情况下没有问题但营销活动那天的爆量流量里混杂了不少爬虫和机器人账号导致状态量疯狂膨胀RocksDB的占用空间整整涨了一倍。当时有两件事必须做一是把状态TTL从两周调整到72小时因为风控场景的特征状态只对最近几天有意义太老的状态其实没价值二是给状态增加清理策略Flink的RocksDB状态支持增量清理和全量清理两种模式我们在低峰期触发一次全量清理把无效状态一次性清了出去。状态瘦身后作业的恢复时间也降下来了原来从Checkpoint恢复要五分钟现在不到一分钟就能拉起。6.4 一个容易被忽略的坑内存参数与Checkpoint频率还有一次线上事故是Checkpoint一直做不完导致整个作业频繁重启。排查下来发现两个问题一是Checkpoint的间隔设得太激进1秒一次而每次快照的数据量比较大反而拖垮了性能二是TaskManager的内存分配不合理托管内存给了256MB实际需要1GB多导致状态写入RocksDB的时候频繁触发磁盘写队。这些参数不是套个默认值就完事的。我的习惯是先按每个算子最大预期状态量 30%冗余分配托管的堆外内存再根据业务可接受的故障恢复时间 数据积压量 / 消费速率反推Checkpoint间隔。比如业务要求最多容忍1分钟的数据回放那Checkpoint间隔就设为30秒左右留一半的余量给恢复时间。合理配置下的实时决策管线容错能力会完全不一样。最后分享一个经过多次踩坑后沉淀下来的心得流处理做实时决策最忌讳的就是重计算、轻工程。数据接入的格式、内存的分配、状态的TTL、Watermark的容忍度这些写不进论文但在生产环境里天天咬人的细节才是决定一个AI原生应用能不能真正支撑实时决策的关键。先把这些基础工程的东西做扎实了再去谈模型的优化路才会走得稳。