简介:这份PDF文档围绕基于Flink的在线机器学习系统架构展开,面向大数据与机器学习方向的工程师、架构师及技术研究者,帮助读者理解如何借助Flink的流批一体能力实现机器学习实时化。文档共1个PDF文件,压缩包约2.91MB,内容以架构图、流程示意与关键技术讲解为主,便于快速通读与查阅。目前已有248人学习下载。文档系统梳理了实时机器学习系统的完整工作流,涵盖数据处理、特征工程、模型训练、模型更新与模型部署五个阶段,并深入介绍Flink流式处理与批处理能力、AI Flow统一训练验证部署流程、事件驱动调度机制以及Flink AI Flow整体架构等核心知识点。读者可从中获得从离线样本到实时样本、从静态特征到动态特征、从T+1更新到增量训练的演进思路,以及流批统一训练与在线推理服务对接的架构设计参考,适合用于技术选型、方案设计与团队内部分享。
1. 在线机器学习为什么要和 Flink 绑在一起
离线训练、定时批推的模式,很多团队都跑过:T+1 跑一遍特征、训一版模型、第二天上线。问题是业务等不起。风控要在一笔交易发生的几百毫秒内判断风险,推荐要在用户滑动的间隙调整排序,广告出价要在一次竞价窗口内完成预估。这些场景里,模型必须在线更新、在线推理,数据一到就得算,算完就得用。
在线机器学习系统架构要解决的核心矛盾就三个:数据流是无限的、模型是要持续迭代的、服务是要低延迟稳定的。Flink 之所以常被选作这套架构的底座,是因为它天生就是为无界流设计的,带状态、带事件时间、带 Exactly-Once 语义,还能把批和流统一到一套 API 里。把特征计算、样本拼接、模型更新、在线推理这几段串起来,Flink 承担的是那条贯穿始终的数据主干。
这篇面向的是准备把在线机器学习真正落地的一线工程师:你可能已经会用 Flink 写作业,但不确定特征和模型该怎么接;也可能模型服务已经跑起来了,但特征和训练对不上。下面按「架构怎么分层 → 特征和样本怎么在 Flink 里做 → 模型怎么更新和推理 → 坑在哪 → 怎么验证」推一遍,能照着搭出最小可跑版本。
2. 在线机器学习系统的分层架构与 Flink 的定位
2.1 四层结构:数据、特征、模型、服务
一套能跑的在线机器学习系统,我一般拆成四层,每层职责清晰,别混在一起。
数据层负责原始事件的接入和缓冲。常见做法是业务埋点或 CDC 变更日志先进消息队列,Flink 作业从队列消费。这一层要保证的是顺序和可重放,别在这里做任何业务计算。
特征层是 Flink 的主战场。它做两件事:一是实时特征计算,比如过去 5 分钟某用户的交易笔数、过去 1 小时某商品的点击率;二是特征拼接,把实时算出来的特征和特征存储里查到的离线特征拼成一条完整样本。这一层的输出有两个下游:写进在线特征存储供推理用,写进样本流供训练用。
模型层负责训练和更新。在线机器学习不等于全都在线训练,常见的是「在线推理 + 准在线更新」:用 Flink 产出的样本流做增量训练或触发全量重训,训练完把模型推到线上。真正的在线学习(每来一条样本就更新一次)只在少数场景用,因为稳定性和可复现性很难保证。
服务层负责推理。模型加载好之后对外提供预测接口,推理时要拿到和训练时一致的特征。这一层的关键是特征一致性,后面会专门讲。
Flink 横跨数据层和特征层,同时向模型层输出样本。它不负责推理本身,但推理要用的特征几乎都从它这里出。这个定位决定了 Flink 作业的稳定性直接决定整条链路的可用性。
2.2 为什么是 Flink,而不是 Spark Streaming 或自己写消费者
选型上被问得最多的就是这句。我的判断标准是三条:状态管理、事件时间、Exactly-Once。
Spark Streaming 的微批模型在秒级延迟上够用,但它的状态是藏在 RDD 里的,做大规模 keyed state 和定时器(比如「用户 30 分钟没动作就清理状态」)很别扭。Flink 的 KeyedState 和 Timer 是一等公民,特征计算里大量用到「按用户分组、按时间窗口聚合、超时清理」,这套原语用起来顺手得多。
事件时间这块,在线场景里数据乱序是常态。用户的操作日志可能因为网络延迟晚到几十秒,如果用处理时间做窗口,算出来的「过去 5 分钟交易笔数」是错的。Flink 的 Watermark 机制能按事件时间推进窗口,配合 allowedLateness 处理迟到数据,这是自己写消费者很难做对的。
Exactly-Once 决定了样本和特征会不会重复或丢失。训练样本一旦重复,模型会被带偏;特征一旦丢失,推理时就会拿到空值。Flink 的 Checkpoint 配合两阶段提交 Sink,能保证从 Source 到 Sink 的一致性。自己写消费者要做到这点,得手动维护 offset 和幂等写入,工作量不小。
提示:如果业务延迟容忍度在分钟级以上,且状态逻辑简单,Spark Streaming 或直接用 Kafka Streams 也能做。Flink 的优势在复杂状态和低延迟同时要的时候才明显。
2.3 最小可跑架构的组件清单
落地时不用一上来就上全套。下面这张表是我搭最小版本时会准备的组件,按优先级排。
| 组件 | 作用 | 最小替代方案 |
|---|---|---|
| 消息队列 | 原始事件接入 | Kafka 单节点 |
| Flink 集群 | 特征计算与样本拼接 | 本地 Standalone 或 MiniCluster |
| 特征存储 | 在线特征读写 | Redis |
| 样本存储 | 训练样本落盘 | 对象存储或 HDFS |
| 模型服务 | 推理接口 | 单进程 Flask/FastAPI |
| 模型仓库 | 模型版本管理 | 本地目录或对象存储 |
这套跑通之后,再考虑把特征存储换成带版本管理的方案、把模型服务做成多副本。别一开始就追求大而全,先把「事件进 → 特征出 → 样本落 → 模型推」这条线打通。
3. 用 Flink 做实时特征计算与样本拼接
3.1 特征计算的三种典型模式
实时特征按计算方式分三类,写法差别很大。
第一类是滑动窗口聚合,比如「过去 5 分钟交易笔数」。用keyBy加slidingProcessingTimeWindow或事件时间窗口都行,关键是窗口大小和滑动步长的选择。窗口越大状态越大,滑动步长越小计算越频繁。
第二类是会话窗口,比如「用户一次会话内的行为序列」。用EventTimeSessionWindows.withGap,gap 设成业务上认为「会话结束」的静默时长。
第三类是自定义状态,比如「用户最近一次登录距今时长」。这种没有固定窗口,得用ValueState或MapState自己维护,配合 Timer 做清理。
// 滑动窗口统计过去5分钟交易笔数,事件时间语义 DataStream<Transaction> transactions = env .addSource(new FlinkKafkaConsumer<>("txn", new TransactionSchema(), props)) .assignTimestampsAndWatermarks( WatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, ts) -> event.getEventTime()) ); DataStream<Feature> txnCount = transactions .keyBy(Transaction::getUserId) .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.seconds(30))) .aggregate(new CountAggregator(), new CountWindowFunction());这段代码里三个参数最关键。forBoundedOutOfOrderness(Duration.ofSeconds(10))表示允许数据迟到 10 秒,设太小会丢迟到数据,设太大会增加窗口触发延迟。SlidingEventTimeWindows.of(Time.minutes(5), Time.seconds(30))表示窗口长 5 分钟、每 30 秒滑动一次,也就是每 30 秒输出一次「过去 5 分钟」的结果。keyBy的字段决定了状态怎么分布,用户量大的话要确认 key 的基数不会导致单点热点。
3.2 特征拼接:实时特征和离线特征怎么合
推理时要用的特征往往一部分是实时的(刚才算的),一部分是离线的(用户画像、商品类目)。拼接的常见做法是异步 IO 查特征存储。
// 异步查询Redis补齐离线特征,避免同步IO阻塞 DataStream<EnrichedFeature> enriched = realtimeFeature .keyBy(Feature::getUserId) .process(new AsyncEnrichFunction()); public class AsyncEnrichFunction extends KeyedProcessFunction<String, Feature, EnrichedFeature> { private transient RedisClient redis; @Override public void open(Configuration params) { redis = new RedisClient("redis-host", 6379); } @Override public void processElement(Feature f, Context ctx, Collector<EnrichedFeature> out) throws Exception { // 异步发起查询,回调里做拼接 redis.get(f.getUserId(), profile -> { EnrichedFeature ef = new EnrichedFeature(f, profile); out.collect(ef); }); } }异步 IO 的核心是别让查询阻塞主流程。Flink 的AsyncDataStream或KeyedProcessFunction里手动做异步回调都行,但要注意两点:一是超时设置,Redis 查不到或超时要给默认值,不能让整条流卡住;二是容量控制,异步请求堆积太多会 OOM,用AsyncDataStream.unorderedWait时设好capacity参数。
3.3 样本拼接:把特征和标签对齐
训练样本需要特征和标签。标签通常是事后才有的,比如「这笔交易是不是欺诈」要等人工审核或用户投诉才知道。所以样本拼接本质是一个「等标签」的过程。
常见做法是把特征先写进一个带 TTL 的状态或外部存储,标签到达时按 ID 回查特征,拼成样本。用 Flink 的KeyedCoProcessFunction可以同时接特征流和标签流,按 key 对齐。
// 特征流和标签流按交易ID对齐,拼成训练样本 DataStream<Sample> samples = featureStream .connect(labelStream) .keyBy(Feature::getTxnId, Label::getTxnId) .process(new SampleJoinFunction()); public class SampleJoinFunction extends KeyedCoProcessFunction<String, Feature, Label, Sample> { private ValueState<Feature> featureState; @Override public void open(Configuration params) { // 特征保留2小时,等标签到达 StateTtlConfig ttl = StateTtlConfig.newBuilder(Time.hours(2)).build(); ValueStateDescriptor<Feature> desc = new ValueStateDescriptor<>("feature", Feature.class); desc.enableTimeToLive(ttl); featureState = getRuntimeContext().getState(desc); } @Override public void processElement1(Feature f, Context ctx, Collector<Sample> out) throws Exception { featureState.update(f); } @Override public void processElement2(Label l, Context ctx, Collector<Sample> out) throws Exception { Feature f = featureState.value(); if (f != null) { out.collect(new Sample(f, l)); featureState.clear(); } } }TTL 设成 2 小时是个经验值,取决于标签到达的最大延迟。设太短会丢样本,设太长状态会膨胀。上线前先统计一下标签延迟的分布,取 P99 再加一点余量。
注意:样本拼接里最容易翻车的是特征和标签的时间对齐。特征必须是标签发生「之前」的,否则就是标签泄漏,训练出来的模型离线指标很好看,上线就崩。
4. 模型更新与在线推理的衔接方式
4.1 三种更新策略:全量重训、增量训练、在线学习
模型怎么更新,直接决定架构复杂度。
全量重训最简单:定时(比如每小时)用最近一段时间的样本重新训一版,训练完推到线上。优点是稳定、可复现,缺点是更新有延迟,且每次训练成本高。
增量训练用新样本在旧模型基础上继续训,更新频率可以更高。Flink 产出的样本流可以直接喂给训练任务。难点是增量训练容易灾难性遗忘,需要混入一部分历史样本。
在线学习是每来一条样本就更新一次模型,延迟最低,但工程上最难。模型要支持在线更新、要能回滚、要防样本噪声把模型带偏。只有少数对延迟极度敏感的场景值得这么做。
我的建议是先用全量重训跑通链路,再根据业务对更新延迟的要求决定要不要上增量。在线学习留到最后,别一上来就啃。
4.2 模型热加载:不重启服务换模型
推理服务换模型不能停服务。常见做法是模型文件放对象存储或模型仓库,服务定时拉取或监听变更,加载到内存后原子替换。
# 模型热加载:后台线程定时检查新版本,原子替换 import threading, time, pickle class ModelServer: def __init__(self, model_path): self.model_path = model_path self.model = self._load(model_path) self.lock = threading.Lock() threading.Thread(target=self._watch, daemon=True).start() def _load(self, path): with open(path, "rb") as f: return pickle.load(f) def _watch(self): last_mtime = 0 while True: mtime = os.path.getmtime(self.model_path) if mtime > last_mtime: new_model = self._load(self.model_path) with self.lock: self.model = new_model last_mtime = mtime time.sleep(10) def predict(self, features): with self.lock: return self.model.predict(features)这里用文件 mtime 做变更检测是最土但最可靠的方式。生产上更常见的是模型仓库提供版本接口,服务轮询版本号。关键是替换时加锁,保证推理请求要么用旧模型要么用新模型,不会读到半个模型。加载失败要保留旧模型,别把服务搞挂。
4.3 特征一致性:训练和推理必须用同一套逻辑
这是在线机器学习里最隐蔽的坑。训练时特征是用 Flink 算的,推理时特征可能是用另一套代码算的,两边逻辑一旦有细微差别,模型效果就会掉。
解决办法是特征计算逻辑只写一份。常见做法是把特征计算封装成独立的库或 UDF,Flink 作业和推理服务都调它。如果推理服务是 Python、Flink 是 Java,那就得保证两边实现严格对齐,或者干脆让推理也走 Flink 的算子。
另一个办法是推理时不重算,直接查在线特征存储。Flink 算好的特征写进 Redis,推理服务按 key 查。这样特征只有一份来源,一致性有保证。代价是推理多一次网络查询,延迟会增加几毫秒。
提示:上线前一定要做特征一致性校验。用同一批原始数据,分别走训练特征链路和推理特征链路,比对输出。差异超过阈值就别上线。
5. 在线机器学习系统架构的避坑与排查
5.1 状态膨胀导致 Checkpoint 越来越慢
现象:作业跑几天后 Checkpoint 时间从几秒涨到几分钟,甚至超时失败。
原因:KeyedState 没有设置 TTL,或者 TTL 设得太长。用户维度的状态随着用户数增长无限膨胀,每个 Checkpoint 都要把这些状态快照出去。
解决:给所有状态加 TTL,按业务需要设最短合理时长。用StateTtlConfig配置,并开启cleanupInRocksDBCompactFilter让 RocksDB 在压缩时清理过期状态。同时确认 key 的基数,如果某个 key 特别热,考虑加盐打散。
5.2 特征和标签时间错位造成标签泄漏
现象:离线评估 AUC 0.95,上线后效果和随机差不多。
原因:样本拼接时用了标签发生之后的特征。比如预测「这笔交易是否欺诈」,却把「交易后 1 小时内是否被投诉」也算进特征了。
解决:拼接时严格按事件时间对齐,特征的时间戳必须早于标签的时间戳。在KeyedCoProcessFunction里加时间戳校验,不满足的样本直接丢弃。上线前做一次特征时间戳的分布检查。
5.3 异步 IO 容量打满导致背压
现象:作业吞吐上不去,Web UI 显示背压,异步 IO 的队列一直满。
原因:AsyncDataStream的 capacity 设得太大,或者下游 Redis 响应变慢,异步请求堆积。
解决:capacity 按「单并行度能承受的在途请求数」设,一般几百到一千。给异步请求设超时,超时的走降级逻辑返回默认特征。监控 Redis 的 P99 延迟,延迟涨了要告警。
5.4 模型热加载时读到不完整文件
现象:换模型后推理报错,或者预测结果异常。
原因:模型文件还在写入时就被加载了,读到了半个文件。
解决:模型文件先写临时文件,写完再原子重命名。加载方只读最终文件名。或者用模型仓库的版本机制,版本号变了才加载,加载前校验文件完整性(比如 checksum)。
5.5 Watermark 停滞导致窗口不触发
现象:特征一直不输出,窗口结果迟迟不来。
原因:某个分区没有数据,Watermark 无法推进。Flink 的 Watermark 取所有分区的最小值,一个分区空闲就会拖住整个作业。
解决:开启withIdleness,让空闲分区不参与 Watermark 计算。WatermarkStrategy.forBoundedOutOfOrderness(...).withIdleness(Duration.ofMinutes(1))。同时检查数据源是否有分区长期无数据。
6. 用离线回放验证在线链路是否真的对
链路搭完,怎么确认它是对的?我的习惯是做一次离线回放:把历史真实数据按事件时间重新灌进 Flink 作业,看产出的特征和样本是否符合预期。
具体做法是准备一份带时间戳的历史事件,用 Kafka 的 producer 按原始时间间隔(或加速)重放。Flink 作业用事件时间语义消费,产出的特征写到一个临时存储。然后拿这批特征和离线数仓里用 SQL 算出来的同口径特征做比对。
比对时重点看三个指标:一是覆盖率,多少比例的事件产出了特征;二是数值一致性,相同 key 相同时间窗口下两边数值的差异;三是延迟分布,从事件时间到特征产出的时间差。覆盖率低说明有数据被过滤或丢失,数值不一致说明计算逻辑有偏差,延迟分布异常说明 Watermark 或窗口配置有问题。
# 用kafka-console-producer按时间戳重放历史数据 # 数据格式:event_time|user_id|amount cat history_events.txt | while IFS='|' read ts uid amt; do echo "{\"event_time\":$ts,\"user_id\":\"$uid\",\"amount\":$amt}" sleep 0.01 # 控制重放速度,0.01秒一条约100TPS done | kafka-console-producer --topic txn --bootstrap-server localhost:9092重放速度用 sleep 控制,想快速验证就设小一点,想模拟真实延迟就按原始间隔。重放期间盯着 Flink 的 Checkpoint 和背压指标,如果重放都扛不住,线上流量来了更扛不住。
验证通过之后,把这次回放用的数据集和比对脚本存下来。以后每次改特征逻辑或升级 Flink 版本,都跑一遍回归,比拍脑袋上线靠谱得多。我自己就吃过亏:有次改了个窗口的滑动步长,觉得影响不大直接上了,结果线上特征分布整个偏了,模型效果掉了好几个点,回滚加排查折腾了一晚上。从那以后,任何特征逻辑的改动都必须过一遍回放验证,这成了我的固定习惯。希望帮到你。
本文还有配套的精品资源,点击获取