ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink流批一体架构解析:从SQL统一到实时机器学习实践

Flink流批一体架构解析:从SQL统一到实时机器学习实践 简介一份聚焦 Flink 流批一体技术架构的解决方案型 PPT面向大数据架构师、实时计算开发人员及技术决策者帮助快速理解流批融合的核心理念与落地路径。内容围绕“技术创新变革未来”展开系统介绍 Flink 的整体架构组件与 Standalone、YARN、K8S 等部署方式重点剖析 DataStream API、DataSet API 与 SQL 统一入口的设计思路并通过在线机器学习平台案例说明高吞吐批处理与低延迟流计算的大规模实践方法同时对比 Lambda 与 Kappa 架构梳理流批一体演进逻辑为技术选型提供参考。内容包含需求和挑战、Flink 架构简介、流批一体入口 SQL、大规模实践及总结展望等模块并附有 Word Count 示例加深对 DataStream API 的理解。资源为单个 PPTX 演示文稿压缩包大小仅 1.13MB结构完整、便于查阅。已有 581 人学习适合需要系统梳理 Flink 流批一体知识体系、规划实时数仓或平台架构的读者收藏使用。1. 流批一体为什么难Lambda 和 Kappa 都没解决的问题流批一体这个概念被讨论了快十年真正落地的并不多。原因在于流和批从计算模型到执行语义都不同批处理追求高吞吐、精确一次、确定性结果流处理要求低延迟、持续输出、能处理乱序数据。早期主流方案是 Lambda 架构用两套引擎分别跑批和流再在服务层合并结果代价是两套代码、两套运维、两套口径数据对不上时很难排查。Kappa 架构试图用纯流解决但 Kafka 这类消息队列在数据回溯和批量 Scan 场景下性能远不如列存系统回溯几个小时前的数据要重新消费很长时间。Flink 的做法是从架构层面统一两者而不是在应用层做适配。这套设计的核心在于三个点用同一个 DAG 描述流批作业、用 SQL 作为统一入口、Runtime 层统一成 push-based 流式执行。底层的统一才是关键否则只是又造了一个 Lambda。本文基于 Flink 流批一体的技术架构 PPT 展开适合正在做实时数仓选型、或者想理解 Flink 内部设计逻辑的工程师。2. Streaming Dataflow 抽象与 DataStream/DataSet API 的分裂期2.1 点边模型Flink 计算模型的最小表示Flink 对作业逻辑的抽象非常简洁——DAG由点和边构成。点是算子operator承载 flatMap、aggregate、keyBy 这类计算逻辑边是数据流通管道可以跑在网络、文件、内存三种介质上。这个抽象在 Flink 0.9 引入流式执行引擎时就已确立但当时批和流用的是两套 API 表达。PPT 里的 Word Count 示例比较有代表性val lines: DataStream[String] env.readFromQueue(address) val words: DataStream[Word] lines.flatMap((line) split(line)) val counts: DataStream[Int] words.keyBy(word).sum(frequency) counts.addSink(new RollingSink(path))这段代码用的是 DataStream API每个转换对应一个 Stream Operator。readFromQueue从消息队列读取无界流flatMap做分词keyBy按词分组sum做频次累加最后写入 RollingSink。关键的语义点是这里的每条数据到达后立即被处理结果持续输出没有“等到所有数据到齐”的概念。对比批处理的 DataSet API同样一个 Word Count 会有本质区别——批处理会先做完 Stage 划分每个算子在读取全部输入后才触发下一阶段算子间通过文件落盘传递数据。这就是当时流批分裂的根源同一个业务逻辑用两套 API 写两遍执行语义还不一致。2.2 无界流与有界流批是流的特例PPT 里有一句话值得注意“Word Count批处理是流计算的特例”。这背后是 Flink 在架构层面做出的关键判断——有界数据流Bounded Stream只是无界数据流Unbounded Stream在“数据全集已知”条件下的特殊情况。这个概念落到执行引擎上意味着不需要两套执行器。流作业是一条从 source 到 sink 的长流水线每条记录逐级穿过算子批作业同样可以表达为这样的流水线只是在边界条件上做区分——source 读取有界数据后发完 Finite 信号触发下游算子的最终状态输出。这意味着 Flink 不需要像 Spark 那样严格区分 Streaming 和 Batch 两套计划执行器而是可以用同一套任务调度框架承载两种模式。关键差异只在几个内部行为上状态后端是否允许落盘、何时触发 Checkpoint、源是否持续监听。这套认知在后续版本中逐渐沉淀最终演变成 1.14 之后的 unified pipeline 设计。2.3 旧架构的分裂代价PPT 明确指出旧架构的问题DataSet API 和 DataStream API 各自有独立的执行路径——Batch Plan → Optimized Plan → Job Graph → Batch Task Driver而流处理走的是Transformation → StreamGraph → JobGraph → Stream Task Operator。两条路径在 JobManager 内部不共享 Planner 和执行计划优化策略。开发者层面体会更深DataSet API 做 join 时可以自动重分区、自动选择 sort-merge join 或 hash join而 DataStream API 的 join 需要手工管理窗口和状态同样的目的两组 API 的调优参数不通用。PPT 里总结了三个痛点语义难以和 SQL 保持一致、添加功能链路长、执行模式不同导致代码无法复用。这三个问题直接催生了后面的架构改造。3. SQL 作为流批统一入口语义一致性与 Early Fire 机制3.1 同一份 SQL两种执行模式PPT 把 SQL 定位为流批一体的入口并给出一个 GROUP BY 聚合的例子USER_SCORES 表包含 User、Score、Time 三列。批模式下跑全量聚合直接对所有历史数据求 sum(Score) 和 max(Time)一条 SQL 在提交时即确定输入数据全集输出一行最终结果。-- Batch Mode SELECT Name, SUM(Score), MAX(Time) FROM USER_SCORES GROUP BY Name;流模式下SQL 语义的关键变化是引入时间维度——因为数据是持续到达的聚合结果会不断更新。这里表格里的示例展示了实时输出时间点Julie (SUM, MAX)Frank (SUM, MAX)12:01(7, 12:01)(3, 12:03)12:03(8, 12:03)(3, 12:03)12:07(12, 12:07)(5, 12:06)同一份 SQL 在流模式下会产生多条中间结果且这些结果随着时间推移被修正。PPT 用窗口区间表达为[-inf, 12:01)、[12:01, 12:04)、[12:04, now)说明这是典型的基于时间进度watermark驱动的窗口计算。3.2 Retraction 机制流上纠正错误结果流模式允许提前输出一部分结果但在事件时间语义下迟到数据会导致之前的结果不准确。Flink SQL 的做法是通过 Retraction 机制修正当结果需要变化时先发送一条标记为 Retraction 的旧值消息再发送一条新的正确值。-- Stream Mode 下同一条 SQL 的底层执行逻辑会附加 Retraction 标记 INSERT INTO result_table SELECT Name, SUM(Score), MAX(Time) FROM USER_SCORES GROUP BY Name;执行时Flink 的 Query Processor 会自动将聚合结果包装成(true/false, row)二元组。false 表示撤回之前的输出true 表示新增或更新。下游算子拿到 false 标记后从关联结果中删除旧记录再 apply 新记录。这就是 PPT 里“流有 Early fire最终结果一致”的本质含义——中间过程不一致终态收敛到和批处理相同的结果。这里可以总结 Retraction 的触发条件基于事件时间窗口的聚合、基于 SQL 的 join特别是维表 join、以及 distinct 类操作。批处理不需要这个机制因为数据全集已知不会出现被修正的中间结果。3.3 Query Processor 模块架构层面的统一入口为了支撑上面这套“同一份 SQL 两种执行模式”的设计Flink 引入了 Query Processor 模块。它位于 Table API SQL 和 Runtime 之间扮演三层角色Logical Plan、Optimizer、Physical Plan、Execution DAG。SQL Table API ↓ Logical Plan — 抽象语法树与执行模式无关 ↓ Optimizer — 基于成本优化 基于规则的优化 ↓ Physical Plan — 流模式映射为 Stream Physical Plan — 批模式映射为 Batch Physical Plan ↓ Execution DAG — 统一到 DAG API Stream Operators这个模块的价值在 PPT 中体现为四条架构改造点Table API/SQL 升级为一级 API引入 Query Processor 统一流批处理路径使用相同的 DAG 和 Stream Operator 描述作业Runtime 统一到流上的 push-based 实现。其中第四点最关键——批处理不再有独立的执行引擎调度而是与流式执行共用同一套调度和容错机制。4. 大规模实践在线机器学习平台的样本生成与数据回溯4.1 场景拆解Event、Entity 与 Sample 三类数据的存储选型PPT 中在线机器学习平台的案例非常具体涉及三类数据。Event 是用户行为事件包括曝光、点击、购买等Entity 是准静态特征比如商品 7 天点击量Sample 是训练样本由 Event 和 Entity 拼接而成。这里第一个要解决的问题是存储选型。实时特征计算需要低延迟的流式订阅历史回溯又需要高吞吐的批量 Scan单一存储很难同时满足。PPT 给出的方案是“消息队列 类 HBase 的 KV 系统”Kafka 提供低延迟流式订阅HBase-like 系统提供点查和范围 Scan。ETL 作业使用 At least once 语义写入 KV 系统利用 KV 的 update 能力完成幂等去重避免重复数据。-- 从消息队列消费事件写入 KV 系统去重逻辑依赖于 KV 的 put(key, value) 覆盖语义 INSERT INTO event_store SELECT event_id, user_id, action, ts FROM kafka_source -- Flink 层不做过重去重仅依赖 At least once -- 下游 KV 的 upsert 能力保证最终一致注意这里的设计选择值得学习没有在 Flink SQL 里硬做 Excatly once 的精确去重。PPT 明确说 Checkpoint barrier 对齐会导致延迟波动因此选择 At least once KV upsert牺牲极小概率的重复读取换取更平稳的延迟。在实际生产中这就是延迟一致性和精确性之间的经典权衡。4.2 CVR 样本生成点击到成交的时间窗口与 Retraction 应对CVR 模型转化率模型的场景非常典型用户点击后是否达成购买。点击和购买之间可能间隔几秒到几小时因此最初生成的负样本点击未成交可能在三小时后因为用户完成购买而变成正样本。PPT 给出的解决方案是 Flink SQL 的 Retraction 机制在时间容忍窗口内先输出基于当前数据的样本当新事件导致结果变化时Broker 先输出标记为 Retraction 的旧样本再输出修正后的新样本。算法端消费时需要对 Retraction 消息做处理——将之前已写入训练样本库的旧样本标记为无效或者直接更新对应 key。-- 实时样本生成核心逻辑join 点击流与成交事件流 CREATE VIEW cvr_samples AS SELECT c.click_id, c.user_id, c.item_id, CASE WHEN o.order_id IS NOT NULL THEN positive ELSE negative END AS label FROM click_stream c LEFT JOIN order_stream o ON c.click_id o.click_id AND c.user_id o.user_id AND o.order_time BETWEEN c.click_time AND c.click_time INTERVAL 3 HOUR;这段 SQL 的关键在于BETWEEN ... AND ...这个区间条件会触发 Flink 的 state 留存——click 记录需要在状态中保留 3 小时来等待匹配的 order。如果成交发生在 3 小时外该样本永远作为负样本输出。这个超时参数需要算法团队根据业务分布来调PPT 里的经验是在延迟容忍程度内先输出再用 Retraction 修正。4.3 批流样本一致性同一套 SQL 复用与 Source 替换实时训练的一个痛点是样本口径和批量生成不一致。PPT 给出的解决思路非常直接直接复用实时样本生成逻辑一样的 SQL一样的 UDF只在 Source 层做替换。平台将 Source 自动替换成 KV 系统中的历史数据执行引擎自动切换为批处理模式这样同一份逻辑生成两份样本特征口径完全对齐。实时模式: Source Kafka - Stream Join - 输出实时样本 批处理模式: Source HBase-like KV Scan - Batch Join - 输出批量样本实现层面常规做法是把表定义通过 SQL DDL 的WITH子句声明的 connector 参数做两层映射。平台解析 SQL 的逻辑计划将kafka_source替换为kv_source同时把 Query 折叠进批处理 Planner。这样做的好处体现在两个方面一是样本口径统一问题从源头被消除二是批处理作业可以提交到混部资源池在低峰期批量回补样本数据不占用实时集群资源。5. 大规模批处理调优JobManager 性能与 Failover 机制在线机器学习平台的批流一体实践还牵扯到一个容易被忽略的维度——大规模批处理作业对调度和容错的要求。批作业的并发度比流作业高得多上游 N 个并发、下游 M 个并发JobManager 上可能管理数百万条执行边。PPT 里提到三个具体的优化方向避免 N×M 级别的内存占用、JobManager Failover 问题 [FLINK-4911]、Region-based Task Failover [FLIP-1/FLINK-4256]。N×M 问题发生在 Shuffle 连接阶段。上游每个 Task 需要将数据分发给下游所有 Task边的数量是 N×M。如果每个 Shuffle 通道都要在 JobManager 中维护状态和指标内存会随着作业规模指数增长。常规的调优手段是开启taskmanager.memory.shuffle.min和max相关配置同时关注jobmanager.memory.heap.size是否足以容纳作业图元数据。Region-based Task Failover 是一个值得深入理解的机制。默认的 Failover 策略是 Restart-all任何一个 Task 出现异常整个作业全部重启而 Region-based 策略只重启受影响的数据流区域Region。它根据 ExecutionGraph 的拓扑结构划分区域——如果失败节点有输入边则向上追溯到所有可能产生数据的区域一起重启如果失败节点没有输入边只重启该节点所在区域。作业拓扑: Source-A - Operator-B - Operator-C - Sink-D 异常发生: Operator-B 容器宕机 默认策略: 重启整个作业 Region策略: 检测到 Operator-B 失败向上追溯 Source-A 所在区域重启 A、B 区域 Operator-C 与 Sink-D 如果依赖 B 的输出也将级联重启。这个优化对大规模批处理的价值体现在恢复粒度上原本跑 2 小时的批作业因为某个节点 OOM 需要整体重跑Region-based 只需要恢复部分上游节点从文件 Shuffle 中间结果开始续跑整体恢复时间可以从小时级降到分钟级。配合 Flink 的execution.attached和文件 Shuffle 机制批作业的容错性能得到质的提升。最后提一个实践中验证过的指标大规模批作业调优时优先关注 JobManager 的堆内存和 GC 日志。如果Full GC频繁先调整jobmanager.memory.heap.size再检查是否启用了jobmanager.memory.process.size的限制。很多表面上看起来是调度性能的问题实际是 JobManager 元数据膨胀导致的 GC 停顿。用jstat -gcutil pid 1000观察老年代使用率如果持续超过 80%就需要考虑为作业图瘦身——减少不必要的算子链拆分或者将 Shuffle 方式改为与数据规模匹配的策略。这些细节点往往比调整作业并行度更有效果。本文还有配套的精品资源点击获取
返回列表