ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink在线机器学习系统架构:实时特征计算、样本拼接与模型热加载落地指南

Flink在线机器学习系统架构:实时特征计算、样本拼接与模型热加载落地指南 简介这份PDF文档围绕基于Flink的在线机器学习系统架构展开面向大数据与机器学习方向的工程师、架构师及技术研究者帮助读者理解如何借助Flink的流批一体能力实现机器学习实时化。文档共1个PDF文件压缩包约2.91MB内容以架构图、流程示意与关键技术讲解为主便于快速通读与查阅。目前已有248人学习下载。文档系统梳理了实时机器学习系统的完整工作流涵盖数据处理、特征工程、模型训练、模型更新与模型部署五个阶段并深入介绍Flink流式处理与批处理能力、AI Flow统一训练验证部署流程、事件驱动调度机制以及Flink AI Flow整体架构等核心知识点。读者可从中获得从离线样本到实时样本、从静态特征到动态特征、从T1更新到增量训练的演进思路以及流批统一训练与在线推理服务对接的架构设计参考适合用于技术选型、方案设计与团队内部分享。1. 在线机器学习为什么要和 Flink 绑在一起离线训练、定时批推的模式很多团队都跑过T1 跑一遍特征、训一版模型、第二天上线。问题是业务等不起。风控要在一笔交易发生的几百毫秒内判断风险推荐要在用户滑动的间隙调整排序广告出价要在一次竞价窗口内完成预估。这些场景里模型必须在线更新、在线推理数据一到就得算算完就得用。在线机器学习系统架构要解决的核心矛盾就三个数据流是无限的、模型是要持续迭代的、服务是要低延迟稳定的。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.withGapgap 设成业务上认为「会话结束」的静默时长。第三类是自定义状态比如「用户最近一次登录距今时长」。这种没有固定窗口得用ValueState或MapState自己维护配合 Timer 做清理。// 滑动窗口统计过去5分钟交易笔数事件时间语义 DataStreamTransaction transactions env .addSource(new FlinkKafkaConsumer(txn, new TransactionSchema(), props)) .assignTimestampsAndWatermarks( WatermarkStrategy.TransactionforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, ts) - event.getEventTime()) ); DataStreamFeature 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阻塞 DataStreamEnrichedFeature enriched realtimeFeature .keyBy(Feature::getUserId) .process(new AsyncEnrichFunction()); public class AsyncEnrichFunction extends KeyedProcessFunctionString, 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, CollectorEnrichedFeature 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对齐拼成训练样本 DataStreamSample samples featureStream .connect(labelStream) .keyBy(Feature::getTxnId, Label::getTxnId) .process(new SampleJoinFunction()); public class SampleJoinFunction extends KeyedCoProcessFunctionString, Feature, Label, Sample { private ValueStateFeature featureState; Override public void open(Configuration params) { // 特征保留2小时等标签到达 StateTtlConfig ttl StateTtlConfig.newBuilder(Time.hours(2)).build(); ValueStateDescriptorFeature desc new ValueStateDescriptor(feature, Feature.class); desc.enableTimeToLive(ttl); featureState getRuntimeContext().getState(desc); } Override public void processElement1(Feature f, Context ctx, CollectorSample out) throws Exception { featureState.update(f); } Override public void processElement2(Label l, Context ctx, CollectorSample 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(targetself._watch, daemonTrue).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 算的推理时特征可能是用另一套代码算的两边逻辑一旦有细微差别模型效果就会掉。解决办法是特征计算逻辑只写一份。常见做法是把特征计算封装成独立的库或 UDFFlink 作业和推理服务都调它。如果推理服务是 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 版本都跑一遍回归比拍脑袋上线靠谱得多。我自己就吃过亏有次改了个窗口的滑动步长觉得影响不大直接上了结果线上特征分布整个偏了模型效果掉了好几个点回滚加排查折腾了一晚上。从那以后任何特征逻辑的改动都必须过一遍回放验证这成了我的固定习惯。希望帮到你。本文还有配套的精品资源点击获取
返回列表