简介:这是一份关于Flink流批一体技术架构的PPT方案,面向大数据架构师、实时计算开发人员及技术决策者,旨在解析Flink如何通过统一架构实现流处理与批处理的融合。PPT共1个文件,大小约1.13MB,内容精炼,便于快速阅读和汇报使用;目前已有583人学习浏览。方案从需求与挑战出发,介绍了Flink的Storage、Local、Cluster、Standalone、YARN、K8S等部署形态,以及DataStream API与DataSet API两大编程接口;重点讲解了以SQL作为流批统一入口的设计思路,并通过Word Count示例展示批量与流式查询的一致性。同时对比了Lambda架构与Kappa架构的优劣,结合在线机器学习平台的大规模实践,阐述了低延迟流计算、高吞吐批处理、编程接口统一、代码复用等核心优势。读者可借此快速建立对Flink流批一体架构的全局认知,并参考其中关于Streaming Dataflow的点和边抽象、算子链路等内容,为自身实时计算平台选型与架构升级提供借鉴。
1. Flink 流批一体架构:从 PPT 到落地的关键解读
这份《Flink 流批一体的技术架构介绍》是我见过少有的能把流批一体讲透的架构资料。它不跟你绕概念,直接抛出结论:批处理是流计算的特例,流批一体不是把两套引擎拼在一起,而是从 Runtime 到 API 全部统一。我在生产环境里维护过 30 多个 Flink 作业,见过太多团队在 Lambda 架构里维护两套代码、两套引擎,最后被数据不一致折腾到崩溃。这份 PPT 里的架构方案,解决的就是这个切肤之痛——一份代码、一样的结果,既支持低延迟流计算,又支持高吞吐批处理。无论你是正在评估流批一体改造的数据架构师,还是被 Lambda 架构双链路折磨的实时计算工程师,这份材料都能直接拿来当技术方案蓝本。
2. 流批一体的核心矛盾:Lambda 和 Kappa 都没解决的统一问题
2.1 主流架构的痛点:为何 Lambda 复杂、Kappa 局限
要理解 Flink 流批一体架构的定位,得先回到主流架构的对比上。PPT 里列出的 Lambda 架构是很多公司的现状:批处理链路用 Spark/Hive 跑 T+1 离线任务,流处理链路用 Flink/Storm 跑秒级实时任务,两条链路各自维护 ETL 逻辑、各自产出结果。这套方案的直接代价是开发效率低——同一个统计口径要在两套代码里各写一遍,而且 SQL 语义在批和流环境下经常出现偏差。
Kappa 架构试图用一套流引擎搞定所有场景,但它有一个硬伤:历史数据回溯能力弱。Kafka 里的消息有留存周期,超过保留时间的数据无法重放,这就导致 Kappa 架构很难支撑需要全量历史数据的场景,比如模型训练样本的批量生成。PPT 里提到流批一体系统需要同时具备低延迟的流式订阅和高吞吐的历史数据回溯能力,这正是两种主流架构各自的盲区。
2.2 Runtime 统一:批处理是流计算的特例
PPT 里的核心论点是「批处理是流计算的特例」。这个论断的支撑点在架构图里很清晰:Flink 把作业表达为 Streaming Dataflow,也就是由点和边构成的 DAG。点是算子(source、flatmap、aggregate、sink),边是数据流通管道,运行在网络、文件、内存等介质上。关键差异在有界和无界:无界流对应流作业,数据持续到达;有界流对应批作业,数据读取完毕作业即结束。
这意味着不用再为批处理和流处理分别设计执行引擎。我之前接到过一个需求:用户的购买行为数据既要做实时转化率统计,又要每天跑一次全量回刷。老方案是流作业一套 Flink SQL、批作业一套 Spark SQL,统计口径对齐花了两周。用 PPT 里的统一架构,直接在 Runtime 层用同一套 DAG 描述两类作业,流作业处理无界流,批作业把 HDFS 上的历史分区当作有界流,算子和执行逻辑完全复用。这个设计把「流批一体」从口号落到了执行引擎的底层实现上。
2.3 新架构的五项关键改动
PPT 里给出了新架构的五项主要修改点,这部分是理解演进方向的钥匙:
第一,Table API 和 SQL 升级为一级 API,不再依附于 DataStream/DataSet API。第二,引入 Query Processor 模块,统一流和批的处理逻辑。第三,使用相同的 DAG 和 Stream Operator 描述流批作业。第四,Runtime 统一到流式 push-based 实现。第五,未来考虑和 DataStream 共享算子。
我重点说下第四点。push-based 意味着数据由上游主动推送给下游,这和批处理时代拉取式的模型有本质区别。拉取式模型天然适合有界数据,但无法处理无界流——你不知道数据什么时候来,所以只能被动等待。统一到 push-based 后,批作业的有界流也可以被看作一种特殊的流数据管道,数据推完作业自然结束。这个改动我理解是 Flink 1.x 之后 Runtime 演进的核心方向,它让 Query Optimizer 生成的物理执行计划可以做到真正的一份计划两种执行。
3. SQL 作为流批一体入口:Batch Mode 和 Stream Mode 的执行差异
3.1 同一份 SQL 的两种执行语义
PPT 里用 USER_SCORES 例子把 SQL 入口讲得非常直观。同一张表,在 Batch Mode 下执行SELECT Name, SUM(Score), MAX(Time) FROM USER_SCORES GROUP BY Name,得到的是确定性结果——数据全部到齐后计算,输出完整聚合。但在 Stream Mode 下,数据是持续到达的,同一句 SQL 会生成多个窗口期的中间结果,PPT 里展示的[-inf, 12:01)、[12:01, 12:04)、[12:04, now)就是时间窗口的切片。
这里要理解两个核心概念:Early fire 和最终结果一致。Early fire 是流模式下的特殊机制——为了低延迟先输出中间结果,后续数据到达后再触发更新。PPT 里明确说「流有 Early fire,最终结果一致」,这解决了之前很多团队对流批 SQL 的误解:不是流模式算错了,而是中间结果带有时间维度语义。
如果只看一眼流模式输出的中间聚合结果,可能会觉得数据有问题。比如 Julie 的分数从 12:01 的 7 分,到 12:03 变成 8 分(加上新到的 1 分),再到 12:07 变成 12 分。这就是 Early fire 的体现,最终 12 分和批模式算出来的一致。这种机制要求在流模式下消费结果的应用,必须理解中间结果可能被后续更新覆盖。
3.2 Query Processor:统一流批处理的关键模块
Query Processor 的设计是这套架构的枢纽。从 PPT 的模块图看,它包含 Logical Plan、Optimizer、Physical Plan、Execution DAG 四个阶段,大部分组件流批共用,只有最终执行方式不同。
实际配置 Flink SQL 作业时,这直接影响开发方式。流作业和批作业的 SQL 逻辑可以完全复用:
-- 批模式执行:读取 HDFS 上的历史分区 SET execution.runtime-mode = batch; CREATE TABLE user_scores ( user_name STRING, score INT, event_time TIMESTAMP(3) ) WITH ( 'connector' = 'filesystem', 'path' = 'hdfs://namenode:8020/data/user_scores/dt=20240101', 'format' = 'parquet' ); SELECT user_name, SUM(score), MAX(event_time) FROM user_scores GROUP BY user_name;这段 SQL 在批模式下会触发完整的查询优化流程——列裁剪、谓词下推、分区剪枝都会生效。TM 资源分配按批作业模式走,读取完数据后作业自动结束。
-- 流模式执行:从 Kafka 消费实时数据 SET execution.runtime-mode = streaming; CREATE TABLE user_scores ( user_name STRING, score INT, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'user_scores', 'properties.bootstrap.servers' = 'kafka-1:9092,kafka-2:9092', 'properties.group.id' = 'user_scores_group', 'format' = 'json', 'scan.startup.mode' = 'latest-offset' ); SELECT user_name, SUM(score), MAX(event_time) FROM user_scores GROUP BY user_name;同一条统计逻辑,流批之间只是 SET 参数和 source 定义不同。这在之前 Lambda 架构里是不可想象的——批处理要写 Hive SQL,流处理要写 Flink SQL,语法有差异,UDF 要各写一套。统一到 Query Processor 后,UDF 也只需要实现一次。
3.3 SQL 作业的参数配置清单
实际跑 SQL 作业时,参数配置直接影响执行表现。基于我在生产环境调试的经验,整理一份常用的参数对照:
| 参数 | 批模式推荐值 | 流模式推荐值 | 配置说明 |
|---|---|---|---|
| execution.runtime-mode | batch | streaming | 作业执行模式,核心开关 |
| pipeline.cpu | 2-4 | 4-8 | 每个 TM slot 的 CPU 核数 |
| taskmanager.memory.process.size | 4-8g | 8-16g | TM 总内存 |
| parallelism.default | 按分区数定 | 按 Kafka 分区数定 | 默认并行度 |
| execution.checkpointing.interval | 不配置 | 30-60s | 流模式必须开 checkpoint |
| execution.checkpointing.mode | 不配置 | EXACTLY_ONCE | 一致性语义 |
| table.exec.state.ttl | 不配置 | 24h 以上 | 状态过期时间,影响数据正确性 |
批模式下不需要开启 checkpoint,因为作业执行完自然结束,失败直接从起点重跑。流模式必须开,否则 TM 崩溃后状态全部丢失。state TTL 是流模式最容易忽略的参数,它控制 keyed state 的过期时间,设置太短会导致状态被清理,设置太长会内存溢出,一般按业务数据迟到上限配置。
3.4 与 DataStream API 的共存策略
PPT 里提到「未来可以考虑和 DataStream 共享算子」,这意味着目前 Table API 和 DataStream API 还是两套算子体系。我在实际项目里的做法是:能用 SQL 表达的统计逻辑一律用 SQL,需要自定义复杂处理逻辑(比如状态编程、窗口内多流 join 的特殊策略)才用 DataStream API。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env); // 注册 Kafka 数据源为动态表 tableEnv.executeSql( "CREATE TABLE clicks (" + " user_id BIGINT," + " item_id BIGINT," + " behavior STRING," + " ts TIMESTAMP(3)," + " WATERMARK FOR ts AS ts - INTERVAL '3' SECOND" + ") WITH (...)" ); // 用 Table API 做流式聚合 Table aggregated = tableEnv.sqlQuery( "SELECT user_id, COUNT(*) AS cnt " + "FROM clicks " + "WHERE behavior = 'click' " + "GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE)" ); // 转成 DataStream 做自定义处理 DataStream<Row> resultStream = tableEnv.toAppendStream(aggregated, Row.class);这种混合模式在线上机器学习平台场景特别实用:常规特征统计用 SQL 搞定,模型个性化处理逻辑用 DataStream 实现,两种 API 在同一个作业里协同工作。需要注意的是,toAppendStream 只适用于 append 模式的查询,如果查询包含聚合操作且结果会更新(比如非窗口聚合),要用 toRetractStream。
4. 在线机器学习平台实践:事件、实体与样本的批流一体处理
4.1 事件与实体的存储选型:消息队列 + KV 系统的组合拳
PPT 里在线机器学习平台的案例是整个架构方案最难的部分,因为涉及数据的一致性、时效性和可回溯性三重约束。平台的核心数据有三类:Event(用户行为事件)、Entity(准静态特征)、Sample(训练样本)。Event 是流式的用户行为数据,Entity 是描述性数据,Sample 是 Event 和 Entity 组合成的训练样条。
PPT 里提到的问题 1「Event 和 Entity 的存储选择」和问题 4「数据可回溯」导向了一个组合方案:消息队列 + 类 HBase 的 KV 系统。消息队列提供低延时流式订阅,KV 系统提供历史数据 Scan 功能。这套选型的原因在于:如果只用 Kafka,历史数据会被过期清理且无法高效扫描;如果只用 HBase,实时流式订阅能力又不够。两者组合后,平台对数据源进行包装,在不同场景下切换存储。
这块我在实践中的理解是:Kafka 负责「现在发生了什么」,KV 系统负责「过去发生了什么」。比如做商品点击转化率(CVR)模型时,用户点击事件实时进入 Kafka,商品 7 天点击量这类准静态特征存入 KV 系统。训练样本需要把两者 join 起来,实时样本用 Kafka + 缓存维表的方式实现,批量样本直接把 HBase 里的历史数据 Scan 出来全量计算。
4.2 数据一致性权衡:At least once + KV 去重方案
PPT 里明确讨论了一个生产环境最棘手的问题:Exactly once 和延迟的权衡。Checkpoint barrier 对齐机制为了实现严格一次语义,会让作业在 barrier 对齐过程中产生延迟波动。尤其在数据倾斜场景下,某个 key 的数据量大,barrier 流动变慢,整体延迟被拉高。At least once 模式没有 barrier 对齐这个包袱,延迟低吞吐高,但数据可能重复。
PPT 给出的解决方案非常务实:ETL 作业使用 At least once,利用 KV 系统的 Update 能力进行去重。这个思路很巧妙——重复数据写入 KV 系统时走 put 操作,同一主键的重复写入被自然覆盖,数据只在最终存储层做一次幂等。消息队列提供流式订阅,KV 系统负责去重和 Scan。
用一段伪代码理解这个方案的实现逻辑:
输入: event_stream (来自 Kafka) 处理: 1. 消费 event_stream,使用 At least once 语义,开启 checkpoint 保证故障恢复 2. 对每条事件提取主键 (如 order_id + user_id + ts) 3. 写入 KV 系统 (如 HBase) 的同一行,put 操作天然幂等 4. 如果需要扫全量历史,直接 Scan HBase 表,无需重新消费 Kafka这个方案在实际集群里解决了一个非常具体的痛点:Kafka 消费端在 rebalance 或任务失败重试时可能重复消费某些分区,如果下游是普通的消息队列,重复消息会进入样本集,造成训练数据分布偏移。有了 KV 系统的主键去重,重复消费问题被彻底屏蔽。
4.3 SQL Retraction 机制修正错误样本
CVR 模型的样本生成有个特殊挑战:用户从点击到成交的时间不确定,秒级到小时级都有可能。如果用户点击后 30 分钟才购买,实时训练作业在点击发生时无法判断是正样本还是负样本——按当前逻辑先输出负样本,成交后又得把负样本修正为正样本。
PPT 里的解决方案是基于 SQL 的 Retraction 机制。这是 Flink SQL 在流模式下的一个重要能力:当聚合结果因新数据而改变时,先发送一条旧结果并标记为 Retraction(撤回),再发送一条新结果。在样本生成场景里,这个机制被用来修正负样本:
-- 实时样本生成作业 INSERT INTO training_samples SELECT user_id, item_id, behavior, CASE WHEN has_purchase THEN 'positive' ELSE 'negative' END AS sample_type, ts FROM user_click_events LEFT JOIN purchase_events ON user_click_events.user_id = purchase_events.user_id AND user_click_events.item_id = purchase_events.item_id;当用户完成购买后,purchase_events 新增一条记录,join 结果变化,Flink 会针对之前的负样本输出一条 Retraction 消息(即删除这条负样本),再输出一条正样本。消费端算法平台处理 Retraction 消息时,把旧的负样本从训练集移除,追加正确的正样本。
这里必须说明一个实际参数调优的问题。Retraction 机制的代价是状态空间膨胀——为了能够撤回旧数据,Flink 需要保留足够的中间状态。实战中我建议把 table.exec.state.ttl 配置得足够大(至少是业务最大延迟时间的 2 倍),否则超期的状态被清理后,旧样本可能无法正确撤回。同时要监控 state 的磁盘占用,Retraction 查询的状态往往比普通聚合大好几倍。
4.4 批量样本生成:SQL 复用 + 存储切换的降维方案
PPT 提到模型还需要定期进行批量样本生成,原因有两层:一是实时训练产生的误样本需要回刷修正,二是部分模型对时效性要求不高,更关注样本准确性,离线批量生成更合适。这里的创新点是批量生成直接复用实时的 SQL 逻辑——一样的 SQL、一样的 UDF,平台自动替换 Source 为 KV 系统的历史数据 Scan,作业自动切换为批处理模式。
这套机制的工程实现要点是平台层的存储切换能力。开发人员写一套 SQL,定义了事件和实体表的读取逻辑,平台在流模式时把事件表指向 Kafka,批模式时指向 HBase 的历史分区。这个能力和 Flink SQL 本身的连接器抽象强相关——连接器隐藏了数据源的物理实现差异,让同一份逻辑可以被不同物理来源执行。
-- 同一份样本生成逻辑,批模式直接扫描 HBase 历史数据 SET execution.runtime-mode = batch; CREATE TABLE user_click_events ( user_id BIGINT, item_id BIGINT, behavior STRING, ts TIMESTAMP(3) ) WITH ( 'connector' = 'hbase-2.0', 'table-name' = 'user_click_events', 'properties.zookeeper.quorum' = 'zk-1:2181,zk-2:2181' ); INSERT INTO training_samples SELECT * FROM user_click_events WHERE dt >= '2024-01-01' AND dt < '2024-01-08';这个作业跑在混部资源上,利用 Hadoop YARN 的空闲资源窗口。我在团队推广这套做法后,批量样本生成效率提升了两倍多,复用代码的比例超过 90%,之前流批各写一套的逻辑只需要维护一份。关键是平台层的 Source 替换机制要设计好,否则每个 SQL 作业都要手动写两套数据源定义,就失去了一体化的意义。
5. 流批一体落地避坑:Checkpoint 延迟、状态内存与 Failover 实战记录
5.1 Checkpoint 屏障对齐导致的延迟抖动
现象:一个 Flink 流作业平时处理延迟在 200ms 以内,某天开始周期性出现 3-5 秒的延迟尖峰。检查监控发现 Checkpoint 时长从正常的 1 秒突增到 8 秒,且 Delay 指标随之上升。
原因:在使用 EXACTLY_ONCE 语义时,数据流中的 Checkpoint Barrier 需要等待所有上游通道的数据都处理完才能对齐。如果某个 key 的数据量突然增大,这个 subtask 的 barrier 移动速度变慢,整个 checkpoint 的完成时间被拉长,处理延迟随之升高。PPT 里提到的「Checkpoint barrier 对齐导致延迟波动」在生产环境非常常见。
解决:第一优先把语义降级为 AT_LEAST_ONCE——在 checkstpoint 配置里设置setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE),去掉了 barrier 对齐,延迟尖峰立即消失。如果业务要求必须 EXACTLY_ONCE,就需要反序开销推进——调整作业拓扑,把数据倾斜的 key 做二次分发,或者单独放宽该算子的并行度。
5.2 大批量状态导致的内存过载
现象:作业运行 3 小时后内存使用率持续增长,GC 越来越频繁,最终 TaskManager 频繁 Full GC 导致作业不稳定。查看 Flink Web UI 的 State Size 指标,发现单个 keyed state 已经达到数 GB。
原因:状态 TTL 没有配置或者配置过长。Retraction 查询和 join 操作都会保留历史状态,数据量持续增长时状态不被清理,最终拖垮内存。这是流批一体实践中最容易被忽视的参数问题。
解决:设置合理的 TTL,按业务数据的最大延迟时间估算——比如点击到成交最多 12 小时,TTL 设置 24 小时。配置table.exec.state.ttl: 24h,让过期状态被自动清理。另外,为状态较多的作业单独设置 TaskManager 内存分配,用taskmanager.memory.managed.size提高状态内存的上限。
5.3 JobManager 故障导致全作业重启
现象:某个作业 7 天内出现两次长时间不可用,每次都是 JobManager 崩溃后所有 TaskManager 断连,作业从最近 checkpoint 恢复,恢复过程耗时 10 分钟以上,期间产生大量实时数据堆积。
原因:PPT 提到 JobManager 需要支持可靠 Failover([FLINK-4911]),但在没有配置 HA 的集群中,JobManager 单点故障意味着整个作业重启。很多团队在开发环境没有配置 HA,直接带到生产环境就踩坑。
解决:生产集群配置 ZooKeeper 或 Kubernetes 高可用模式。在flink-conf.yaml中配置high-availability: zookeeper、high-availability.zookeeper.quorum: zk-1:2181,zk-2:2181,zk-3:2181、high-availability.storageDir: hdfs://namenode:8020/flink/recovery。配合 region-based failover([FLIP-1/FLINK-4256]),让单个 TaskManager 崩溃只重启受影响区域的任务,而不是全部任务。这个优化对降低故障重启时间非常关键——PPT 里的 Failover 优化我实际体验下来,恢复时间从 10 分钟缩短到 1 分钟内。
5.4 Region-based Failover 在数据倾斜下的表现
现象:开启 region-based failover 后,一个算子的单个 subtask 崩溃,恢复过程中业务数据出现侧流——部分窗口聚合没有立即输出,且延迟升高。
原因:region-based failover 只重启上游到故障点的区域任务,如果故障算子上游是数据重分布算子,恢复后该区域的部分状态需要从上游重放,会导致该区域处理延迟与其他区域不一致。数据倾斜场景下,热点 key 所在区域容易触发这种局部恢复。
解决:为高位键做二次随机分发(刷热 key 前缀),把热点分散到多个 subtask。同时确保每个区域的状态不超过单台 TaskManager 的承载能力。如果业务允许,也可以把最大 KEY 的状态 TTL 缩短,减少需要重放的状态量。
5.5 多个表 JOIN 时 Retraction 带来的结果震荡
现象:三表 join 的 SQL 在流模式下输出结果反复横跳——同一条样本一会是正样本,一会又变成负样本,消费端算法团队反馈训练数据抖动剧烈。
原因:多表 join 在流环境下使用 retraction 机制时,任何一个表的更新都会触发整条 join 链路的结果重算。当维度表和事实表的更新节奏不一致时,会出现中间状态反复变化。这本质上是流处理 join 的语义问题,不完全是实现 bug。
解决:评估业务对时效性的真实容忍度,如果不是严格需要秒级,可以直接把这类 join 作业切换为批处理模式,每天定时跑批量样本生成。把实时链路只保留简单独立的统计逻辑,复杂的多实体 join 全部转移到批量链路,避免 retraction 震荡对训练样本质量的影响。
6. 架构验证方法与调优实战:用 SQL 重放机制验证流批结果一致性
流批一体的核心承诺是「一份代码,一样的结果」。但这份承诺要落地,就需要一套验证方法。我的做法是:利用 Flink 的 State Processor API 和批模式重放机制,对流批结果做系统性比对。
具体验证流程是这样的:先用批模式跑一遍 HBase 里的历史数据,产出全量统计结果作为基准。然后启动流模式作业,消费 Kafka 中同时间段的数据,等到所有数据全部消费完毕(在scan.startup.mode设置为earliest-offset的情况下等待窗口期过后),取出最终的聚合快照,和基准结果做 diff。
// 读取批量结果 Dataset<Row> batchResult = spark.read() .parquet("hdfs://namenode:8020/tmp/batch_result") .select("user_id", "sum_score", "max_ts"); // 读取流模式最终状态快照 Configuration conf = new Configuration(); conf.setString("execution.runtime-mode", "batch"); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf); Savepoint savepoint = Savepoint.load(env, "hdfs://namenode:8020/flink/savepoints/savepoint-20240101", new RocksDBStateBackend("hdfs://namenode:8020/flink/rocksdb"));实际验证时,我建议注意三件事。第一,时间窗口要对齐——流模式输出的窗口结果和批模式的时间切分不一定完全一致,需要确保比较口径相同,否则 diff 出来的差异全是假阳性。第二,验证样本要覆盖边界场景——凌晨流量低谷、双十一流量高峰、业务字段新增等极端情况都要单独检查。第三,验证逻辑本身要纳入 CICD 流程——每次修改 SQL 或新增业务逻辑后自动触发一次重放比对,避免人为遗忘。
还有一个实践中的小技巧:流批作业共享的 SQL 逻辑,建议用单独的模块管理,通过 CI 在提交时做语法校验和运行模式双测。我之前遇到过一次事故——在流 SQL 里加了某个批模式的窗口函数语法,结果流作业提交后直接报错,线上链路中断了半小时。从那以后我每次改 SQL 都强制跑一遍批流双模验证:同一份 SQL 先在本地用批模式读历史数据,再切流模式消费测试 topic,两边输出一致才允许上线。这已经成了团队的标准动作,希望对你也有参考价值。验证框验就写到这里,具体到你的业务场景,欢迎在实践中给我反馈。希望帮到你。
本文还有配套的精品资源,点击获取