ARTICLE DETAIL

资讯详情

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

基于Flink的在线机器学习系统架构设计与实践

基于Flink的在线机器学习系统架构设计与实践 简介基于 Flink 的在线机器学习系统架构探讨是一份面向大数据与实时计算工程师的PDF技术资料围绕实时机器学习系统化展开。文档以流批一体为核心理念阐述 Flink 如何统一流式处理与批处理能力支撑模型从样本生成、特征工程、实时训练到在线更新的完整链路并重点介绍 AI Flow 工作流对训练、验证、部署流程的自动化整合。内容包含机器学习各阶段实时化的对比图解、离线 T1 更新向增量实时训练的演进、Lambda 与 Phi 架构分析、事件驱动调度机制及 Flink AI Flow 架构组成可帮助读者理解在线机器学习系统的落地路径与关键技术取舍。资源为单个 PDF 文档压缩包约 2.91MB目前已有 248 人学习下载。适合正在建设实时计算平台、探索实时特征与模型更新机制的工程师作为架构参考。1. 在线学习为什么绕不开 Flink这个标题到底在探讨什么搜这个标题的人手上多半已经有一套跑批的离线模型或者一条用 Python 脚本拉数据、训练、导模型的手工流水线。这套东西最痛的不是模型精度而是样本产生到特征落库的时间差——在线学习要的是秒级甚至毫秒级的样本反馈。基于 Flink 的在线机器学习系统架构就是围绕流式数据进来、特征实时算出来、模型持续更新、服务端立刻用上新版本这四件事搭的一条闭环管道。这篇探讨真正面向的读者是已经在用 Flink 做实时数仓、准备把模型训练和在线推理也搬到流上的人。适合你的落地场景是推荐、风控、实时定价、广告出价这类特征时效敏感的业务如果你只是做离线 T1 报表这套架构对你就是过度设计。2. 为什么选 Flink 做在线学习管道它不是唯一的流引擎但是最合适的骨架2.1 在线机器学习需要什么样的数据管道在线学习和传统离线训练最本质的区别是训练样本不是一次性给全的。用户点击、交易流水、设备上报这些事件以流的形式持续到达你没法等数据齐了再开始算特征因为数据永远不会齐。所以在线学习管道的第一要求是低延迟的事件接入第二要求是对迟到数据的容忍和修正第三要求是状态——同一个用户的特征历史、同一个模型的版本号都属于需要跨时间保存的状态。Flink 的能力模型正好对应这三点事件驱动、基于时间语义的窗口、以及可容错的状态后端。相比之下把离线 Spark 任务改成每五分钟调度一次表面上也在实时但调度粒度、状态管理和事件时间语义都会成为瓶颈。Spark Streaming 的微批模型在数据到达和实际计算之间存在固定间隔窗口边界是批边界而非事件时间边界。对于点击率预估这种对特征新鲜度敏感的模型一个五分钟的微批延迟可能直接影响线上收入而 Flink 的毫秒级处理延迟配合事件时间窗口能把事件发生和特征生效的差距压缩到真实可用的水平。在线学习的另一个隐藏需求是样本的可回溯性。离线训练你可以反复读历史数据在线学习里样本一旦流过窗口就不可再生。管道设计必须在一开始就把原始事件、特征版本、模型版本三者绑定存储否则后续你想重演一次训练过程会发现特征算不回去。Flink 的 Kafka 连接器和 Checkpoint 机制在这方面有天然优势原始事件保留在 Kafka特征计算逻辑保留在作业代码里重演只是换一个提交位点的事。2.2 Flink 的窗口、状态与背压机制对特征计算的意义在线特征的常见计算是过去 N 分钟内的行为聚合。Flink 的事件时间窗口让迟到数据也能按业务时间归入正确的窗口而不是按到达时间硬切。这对午夜后的交易风控尤其重要——用户在北京时间 00:00:30 的点击可能因为网络抖动在 00:01:10 才到达 Kafka。如果按处理时间来窗口这笔样本会被算进错误的时间桶模型学到的规律就偏移了。状态后端是我在实际项目中最看重的能力。特征不只是窗口里的聚合值还包括跨窗口的实体画像。比如该用户历史 7 天累计购买金额最近一次点击距今多少秒这类特征用纯 SQL 的 group by 表达非常别扭但用 Flink 的 Keyed State 写很自然。ValueState 存最近一次事件时间戳MapState 存品类到计数的映射RocksDB 状态后端让单作业可以扛住十亿级别的 key。在线学习里特征数量动辄几百维每个 key 的状态体积是 KB 级选型时要把状态后端的内存预算算进去否则堆内存 OOM 只是时间问题。背压机制的作用常被低估。当模型训练侧消费跟不上时Flink 会让数据在源头减速而不是在内存里堆积到 OOM。在线学习管道里训练侧的消费速度经常因为模型迭代而抖动——一次新的超参搜索可能让训练任务变慢如果管道没有背压传导Kafka 的 Lag 会一路涨到消费端崩溃。Flink 的背压是自动的但你需要监控它持续背压说明下游能力不足临时背压可能是 checkpoint 期间的正常现象两者要分清。2.3 Flink 与 Spark Streaming、Kafka Streams 的取舍对比维度FlinkSpark StreamingKafka Streams处理模型真流式微批真流式状态管理Keyed State RocksDB有限支持依赖外部存储基于 Kafka changelog端到端延迟毫秒级秒级毫秒级精确一次保障成熟配合 Kafka 两阶段提交支持但配置复杂依赖 Kafka 事务连接器生态丰富Kafka/Hudi/JDBC 等丰富与 Kafka 深度绑定在线学习适配度高中中低我的选型经验是如果团队已经有 Kafka 和 Flink 实时数仓的基建就直接在 Flink 上做省的是一套引擎的运维成本而不是代码量。如果场景只是单一 Kafka 主题的轻量处理Kafka Streams 更轻部署不依赖额外集群。Spark Streaming 在在线学习里的价值主要在批量样本的周期重算比如每天凌晨用全量数据修正特征基线而在实时特征计算这条主链路上Flink 的窗口和状态语义确实更贴合需求。Flink 还有一个被反复验证过的优势是生态的自洽Flink SQL 做特征聚合、DataStream API 做自定义样本导出、Flink CDC 做后端数据同步全部跑在同一套集群上运维只需要关注一套引擎的监控和调优。架构探讨落到工程选型最终拼的是一个团队能不能长期运维住这套系统这一点 Flink 的社区活跃度和文档完整度给了团队足够的安全感。3. 在线机器学习架构怎么搭三路分离与最小可运行作业3.1 架构总览样本流、特征流、模型服务三路分离一个能支撑在线学习的 Flink 架构我一般拆成三条互不阻塞的链路。第一路是样本流从 Kafka 读原始行为事件经过清洗、拼接特征、打上模型版本号后落成训练样本这是模型更新的口粮第二路是特征流同一份原始事件同时进入特征计算作业产出用户和物品的实时特征供样本生成与在线推理共用第三路是模型服务不直接跑在 Flink 里而是把 Flink 产出的特征宽表导出到 Redis 或向量数据库模型服务从那里取数推理请求的延迟和训练管道互不影响。三路分离的核心目的不是架构图好看而是故障隔离。样本流作业重启时特征流不该受影响模型服务依赖的 Redis 抖动也不该引发 Flink 作业反压。同时特征流和样本流必须复用同一套特征计算逻辑——这是在线学习里最容易翻车的地方。离线训练脚本里算的特征如果和服务端从 Redis 读的特征口径不一致模型上线后效果必然变差而且这种劣化很难排查。解决方案是把特征计算逻辑收敛到 Flink 作业里训练样本和线上特征都由它产出从机制上消除两套特征代码。3.2 用 Flink SQL 做实时特征加工一个最小可运行作业在 Kafka 已接入的前提下最入门的做法是用 Flink SQL 建动态表然后做窗口聚合。以用户点击行为为例下面这个作业可以原样起一个验证环境跑通-- 读取点击事件 CREATE TABLE click_event ( user_id BIGINT, item_id BIGINT, category_id BIGINT, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic ods_click_log, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id flink-feature-group, scan.startup.mode latest-offset, format json ); -- 每小时聚合产出用户最近1小时点击量与独立品类数 CREATE TABLE feature_hourly ( user_id BIGINT, hour_bucket TIMESTAMP(3), click_cnt BIGINT, category_cnt BIGINT, PRIMARY KEY (user_id, hour_bucket) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://feature-master:3306/feature_store, table-name user_hourly_feature, username flink, password ${FEATURE_DB_PASSWORD} ); INSERT INTO feature_hourly SELECT user_id, TUMBLE_START(event_time, INTERVAL 1 HOUR), COUNT(*) AS click_cnt, COUNT(DISTINCT category_id) AS category_cnt FROM click_event GROUP BY user_id, TUMBLE_START(event_time, INTERVAL 1 HOUR);逻辑说明第一张表定义 Kafka 数据源WATERMARK 声明该流允许事件迟到 5 秒迟到的事件仍按业务时间进原窗口第二张表定义 MySQL 结果表主键被声明为 NOT ENFORCED意思是 Flink 侧不承担唯一性校验这个声明只用来标识更新键让下游知道按什么字段做 upsert。最后一条 INSERT 把流式聚合写出去每一小时每个用户产出一行特征。参数说明里有两个关键点。scan.startup.mode 设为 latest-offset 只适合功能验证生产上建议用 earliest-offset 或指定 timestamp否则作业重启会跳过 Kafka 里堆积的历史消息样本出现断层。GROUP BY 里同时输出 TUMBLE_START是为了给特征打上时间桶标签下游拼接训练样本时需要知道这条特征对应哪个小时窗口否则跨窗口拼接会错位。COUNT(DISTINCT category_id) 在实时场景下是重操作数据量大的时候建议改成 HyperLogLog 近似去重Flink SQL 里对应 APPROX_DISTINCT 的优化项精度损失对模型训练通常可接受。这段作业在生产上第一个会踩的坑就是 JDBC 连接器异常。MySQL 端默认 wait_timeout 是 8 小时Flink 的长连接会被服务端静默断开作业间歇性报 Communications link failure。常见做法是在 WITH 里加 sink.max-retries 3同时把 MySQL 实例的 wait_timeout 调大或者用连接池中间件接管。这个坑我在第 5 章会展开写。3.3 用 DataStream API 自定义 Source 与 Sink模型版本感知与样本导出Flink SQL 能覆盖大部分特征聚合但在线学习里有两类需求 SQL 不太好写。一类是窗口外的跨事件状态特征比如当前用户最近一次点击距今多少秒这类特征依赖精确的状态更新用 Keyed State 更直接。另一类是模型版本感知当新模型训练完、发布到对象存储后你希望特征管道或服务进程尽快切换版本而不是等人工重启任务。自定义 Source 的常见写法是周期性扫描模型版本目录而不是做长连接推送。代码骨架如下public class ModelVersionSource extends RichSourceFunctionString { private volatile boolean running true; private final String versionPath; private String lastVersion; public ModelVersionSource(String versionPath) { this.versionPath versionPath; } Override public void run(SourceContextString ctx) throws Exception { while (running) { String current FileSystemUtils.readLatestVersion(versionPath); if (current ! null !current.equals(lastVersion)) { ctx.collect(current); lastVersion current; } Thread.sleep(5000); } } Override public void cancel() { running false; } }逻辑说明这个 Source 每 5 秒读一次对象存储里的版本文件发现版本号变化就把新版本字符串发射到下游。关键价值在于让模型更新成为一个普通的数据流事件——下游算子可以用 MapState 记录当前生效版本新样本自动带上新版本号即使用户维度的旧状态还在也能区分样本是哪个模型版本产生的。参数说明轮询间隔不宜小于 3 秒否则对象存储的 GET 请求量会被放大版本文件建议用一个递增数字或提交哈希命名不要用时间戳字符串否则同一版本重传时会误判成新版本。自定义 Sink 的典型场景是把训练样本导出到数据湖。Flink 生态里 Hudi 的 upsert 能力对按主键更新的样本表很友好但如果你只需要把样本落成 Parquet 文件供定期训练FileSink 加 RollingPolicy 就够不必引入整套数据湖表格式。FileSink 的常见配置是每个小时切一个文件、64MB 触发滚动小文件问题会直接影响后续训练任务读取 Parquet 的吞吐这个在第 5 章也会提到。3.4 模型更新闭环训练、评估、发布怎么接回 Flink三路分离之后模型更新闭环是这样串起来的Flink 样本流把带特征和版本号的样本写入数据湖 → 训练平台定时或按样本量阈值触发训练 → 训练产出新模型文件并写回对象存储 → 版本目录更新 → Flink 的 ModelVersionSource 感知新版本 → 特征流给新样本打新版本号 → 模型服务加载新模型并切流。评估环节我建议放在发布之前从数据湖里抽最近 24 小时的样本做离线评估指标达标才更新版本文件上线后要盯推理服务的平均响应时延和模型预测分布这两个指标能暴露特征不一致和数据漂移两类问题。这个闭环里最容易断的环节是特征版本和模型版本的对应关系。如果在样本落库时没有记录特征计算时刻的特征版本号后续模型训练时根本无法判断这条样本用的是哪套特征逻辑。我的做法是在样本宽表里加两个字段feature_version 和 model_version前者是特征计算作业的版本号后者是该样本参与训练时对应的模型版本。有了这两个字段训练数据筛选、模型回放、异常样本定位都变得可操作这份架构才不会变成黑匣子。4. 在线学习架构的四个命门状态、反压、容错与一致性的设计取舍4.1 状态管理特征状态与模型版本状态如何共存在线学习作业里同时存在两类状态设计不清晰就会互相干扰。第一类是业务状态比如用户最近一次点击时间、近 7 天活跃天数这类状态变化频繁且直接参与特征计算第二类是元数据状态比如当前生效的模型版本、特征版本号这类状态变化低频但必须精准。我的经验是两类状态不要混在同一算子里业务状态放 Keyed State元数据状态可以用 Broadcast State 或独立的宽表存储否则一次模型升级会触发全量 key 的状态迁移代价极高。RocksDB 和堆内存的选型也要提前定。特征量级在千万级别以下用堆内存状态简单直接超过这个量级RocksDB 是必选项但要为它规划独立的磁盘空间并且接受单 key 读写延迟从微秒级升到毫秒级。状态 TTL 是个容易忽略的参数。用户短期兴趣特征一般设 7 天长期画像设 30 天不设 TTL 的状态会随着时间无限膨胀最终拖垮 Checkpoint 的性能。设了 TTL 之后要注意Flink 对过期状态的清理是惰性的内存压力不会立刻下降需要结合状态后端的 compaction 调优。4.2 反压机制训练侧消费慢时如何自愈在线学习管道里特征作业和样本导出作业的消费速度往往不对等。特征作业简单吞吐高样本导出要写对象存储还要做拼接容易成为瓶颈。当同一份源数据被两个作业消费时一个作业反压不会影响另一个但如果你把特征计算和样本导出放在同一个作业里慢的 Sink 会拖垮整条链路。架构上要把这两个算子拆开中间用 Kafka 解耦特征结果写一份到 Kafka 特征主题样本导出作业再订阅这个主题两边各自维护 Checkpoint。反压监控要看两个指标TaskManager 网络缓冲区的使用率以及 Kafka Lag。缓冲区使用率持续高于 80%说明下游处理速度跟不上需要扩容下游算子并行度Kafka Lag 涨了但缓冲区使用率不高说明消费端故障或者 checkpoint 卡住。这两种现象的处理方式完全不同最容易犯的错是一看到反压就盲目加并行度结果下游是写数据库加了并行度把数据库打垮反压更严重。4.3 Checkpoint 与端到端精确一次样本链路不重不漏在线学习对样本重复的容忍度很低。一条样本被重复计算一次等于在训练集里人为制造了偏差。Flink 的端到端精确一次靠 Checkpoint 加 Kafka 两阶段提交实现数据源记录消费位点Sink 在 Checkpoint 完成时才提交外部系统的写入事务。这套机制能成立的前提是上下游都支持事务或幂等Kafka 支持事务JDBC Sink 只能做到幂等所以样本导出我建议走数据湖的 upsert 或 Kafka 事务而不是直接写 MySQL。Checkpoint 参数是经验值我的默认配置是间隔 60 秒超时 10 分钟MinPauseBetweenCheckpoints 设为 30 秒。间隔太短会让反压频繁间隔太长会让故障恢复时重放的数据量过大。在线学习管道还要注意一个特殊情况样本流的数据是高频的Checkpoint 恢复时从 Kafka 重放的数据可能让下游训练任务瞬间压力飙升。缓解办法是把 checkpoint 的存储路径放在高性能存储上同时给训练任务的消费端做限流避免恢复后雪崩。5. 在线学习架构避坑指南五个真实翻车案例5.1 JDBC 连接器间歇性报错现象、原因、解决现象Flink 作业运行两三个小时后日志开始报 Communications link failure作业反复重启Kafka Lag 持续上涨。原因MySQL 服务端 wait_timeout 默认 8 小时但 Flink JDBC Sink 的长连接如果空闲时间超过某个阈值会被服务端断开Flink 连接器没有自动重连或者重连逻辑不完善。解决在 JDBC WITH 里加 sink.max-retries 3同时把 MySQL 实例的 wait_timeout 调到 24 小时以上或者用 HikariCP 连接池托管。另一个隐藏因素是要检查 MySQL 的 max_allowed_packet批量写入时一条请求超过包大小也会触发连接中断。5.2 Sink 到 Hive 表数据不入表现象用 Flink 写 Hive 表作业正常结束但查 Hive 看不到数据。原因Flink 写 Hive 默认走的是分区目录如果表是静态分区而作业没有显式指定分区值数据会被写进临时目录或错误分区此外 Flink 和 Hive 的元数据缓存不一致作业创建的表结构在 Hive 侧没有同步。解决先确认表是动态分区还是静态分区动态分区在数据量大时是性能陷阱建议改成按时间字段转字符串显式指定分区写完数据后用 Hive 的 MSCK REPAIR TABLE 同步分区元数据。这个坑在实时特征结果落 Hive 做离线分析时经常遇到建议直接走 Hudi 或 Iceberg 规避。5.3 窗口乱序导致样本丢失现象特征聚合结果比离线统计少业务上完全对不上。原因事件时间窗口配了 WATERMARK但 watermark 的延迟设置小于数据真实乱序程度。网络抖动让部分事件的业务时间比 Kafka 到达时间晚了几分钟甚至更久超过 watermark 的 5 秒容差后迟到事件被丢弃。解决先把 watermark 放宽到 30 秒观察数据用 Flink Web UI 的 EventsIn/EventsOut 对比确认丢了多少如果乱序确实严重要考虑用 allowedLateness 配合迟到数据侧输出把超晚事件单独收集做补偿修正而不是直接丢弃。在线学习场景里被丢弃的样本可能导致模型对某些时段的规律完全没有感知这是比离线报表丢失更隐蔽的坑。5.4 模型文件更新了作业却一直用旧版本现象模型版本文件已更新但推理结果没有变化查看作业日志发现 ModelVersionSource 没采集到新版本。原因对象存储的版本文件路径写错或者文件更新被缓存了。更常见的是版本对比逻辑只比较了文件名而文件名是时间戳同一秒内重新上传产生了同名覆盖比较逻辑判断相同。解决版本号不要用时间戳用 git commit 哈希或自增 ID读取对象存储时加 Cache-Control: no-cache绕过网关缓存。另外轮询间隔和 Flink 的 Source 并行度要匹配Source 并行度为 1 时只有第一个 subtask 在轮询如果有多个实例在跑要确保版本号是全局一致的。5.5 反压和 OOM 同时出现先扩容还是先查状态现象作业频繁 OOM同时 Web UI 显示高反压团队第一时间加并行度结果没改善。原因并行度不是瓶颈状态后端内存才是。加并行度会扩大状态分片每个 TaskManager 的 RocksDB 或堆内存反而更紧张。解决先看 GC 日志和状态大小指标确认是状态膨胀还是数据流量增大。状态膨胀就查 TTL 和状态清理配置数据流量增大才增加并行度而且要同时调整 Kafka 分区数和 Source 并行度对齐否则加了也是空转。这条经验值五万块钱加到错误的地方比不加还糟。6. 从架构探讨到线上验证正确性校验与两个进阶方向架构探讨类的方案最大的风险是画图时很美、落地上没法验收。我给自己定的验证清单有三项第一项是回放测试把过去 24 小时的 Kafka 消息按相同顺序重新灌进管道对比实时特征和离线批特征的一致性差异率要低于 1%第二项是延迟监控统计从事件进入 Kafka 到特征写入特征库的 P99 延迟推荐场景应该控制在 5 秒以内第三项是模型指标监控不单看 AUC要看在线推理的预测均值是否异常偏移这是特征不一致的最早信号。这三项跑通架构才算数。进阶方向我认为有两个值得投入。一个是把 Flink CDC 引入样本回流链路线上服务的变更数据通过 CDC 同步到数仓让训练样本能包含事后才知道的结果标签这对点击率预估的延迟反馈建模很有效。另一个是建立特征血缘管理每个特征字段记录它来自哪个 topic、哪段 SQL、哪个版本排查问题时能顺着血缘定位是上游数据变了还是特征代码变了。这两个方向都不会让架构复杂度失控但对在线学习系统的长期可运维性是质变。我做这套系统最大的教训是别让在线变成口号先问自己业务能不能承受两分钟的特征延迟如果不能承受再谈 Flink 集群的规模规划如果能承受也别一上来就上全部组件从一条样本流、一个特征作业、一台测试模型服务开始跑通闭环再谈扩展。用最小的闭环验证架构可行性比把架构图画完整要重要得多。希望这篇探讨能帮你少走几段弯路。本文还有配套的精品资源点击获取
返回列表