ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Flink批流一体核心:有界流与无界流深度解析及选型实践

Flink批流一体核心:有界流与无界流深度解析及选型实践 做了一个多月的Flink实时任务回头发现组里一半的实时需求其实用批处理跑更快另一半离线报表反而应该用流处理来削峰填谷。问题出在哪很多人对Flink里有界流和无界流的理解只停留在字面上——一个有终点一个没终点。但真正落到代码和运维上这两个概念牵扯出来的是调度策略、状态生命周期、容错方式、资源模型的一整套差异。这篇文章不打算从头讲Flink是什么而是直接拆解批流一体这层窗户纸到底什么是有界流和无界流Flink凭什么能用一套API同时处理两种场景以及在实际项目中怎么选型、怎么避坑。适合已经写过几个Flink任务、但还没系统理清底层逻辑的同学。读完你至少能回答三个问题你的数据到底算有界还是无界该用STREAMING还是BATCH模式批流一体的一致性承诺为什么不能简单画等号1. 有界流和无界流到底差在哪别让“数据没有终点”这句话骗了你1.1 源头决定属性而不是“计算代码”决定属性所有Flink程序的核心都是对数据流的转换。即使你读的是一个文件、一张MySQL表、或者一个Parquet目录Flink也会把它们抽象成“流”。区别只是这个流自带一个bounded属性有界流意味着流中的元素总数是已知且有限的处理完最后一个元素后任务会正常结束无界流则意味着数据源会无限地产生数据任务理论上永远不结束。这里的判断标准只有一个数据源头本身是不是有限集合。读取一个固定文件列表比如/data/2024/01/*.parquet是有界流。读取某个数据库表的全量快照是有界流。监听Kafka的某个topic并持续消费最新消息是无界流。监控一个文件目录并自动处理新落地的文件虽然看起来“像是读文件”但源头是不断扩张的所以它是无界流。很多人在这一步就栽了。我之前接过一个需求业务方坚持说自己要的是“批处理”因为数据是日志文件每天凌晨定时同步到HDFS。我一看代码用的是StreamingFileSource配合FileMonitoringFunction目录一直在追加新文件这其实是无界流。任务跑起来确实不会失败但它永远不会结束也不会主动释放资源。真正的批处理应该是FileSource的Bounded模式或者直接用一个BATCH作业读取固定文件列表。所以不要再写“这个任务是批任务所以数据是有界流”这种话。有没有界请先去看数据源再看运行时模式。BATCH模式只是Flink针对有界流做的一系列优化手段的组合它不会把无界数据变成有界数据。1.2 两种流在运行时的表现差异比你想的大得多用一张表可以看得很清楚维度有界流Bounded Stream无界流Unbounded Stream数据读取策略一次性读取或分批拉取读完后source结束持续拉取source从不结束任务终止所有数据处理完作业自动Finish理论上永不停止调度方式可以按阶段调度算完一批再调度下一批所有算子必须同时启动整条流水线持续运行Shuffle策略可以做Blocking Shuffle边落盘减少上游重算需要Pipelined Shuffle流式边实时往下游推数据状态生命周期状态随作业结束销毁不需要长期保持状态长期保留必须依赖Checkpoint和恢复机制Watermark通常读完后发送Long.MAX_VALUE水位线强制触发所有窗口需要周期性生成水位线还要处理空闲分区资源占用可以动态申请阶段资源处理完即可释放常驻资源功耗和成本可视化故障恢复单个任务重启或阶段重放即可必须整体恢复到最近一次Checkpoint实际运维中的观感差异非常强烈。有界流作业像MapReduce启动跑结束释放。无界流作业像一个永续服务需要监控延迟、吞吐、背压、Checkpoint失败率还要处理Kafka消费位点回溯、窗口错峰、连接器异常。1.3 一个容易误判的场景Kafka里的数据“现在没新的”不等于有界我见过不少新手把Kafka topic里暂时没有新消息误判为“数据有限”。这是典型的概念错位。Kafka topic本身是一个无限流抽象——它的生命周期和你业务的写入节奏无关只要生产者还在写它就是无界即使生产者停了从Flink的角度看这仍然是一个无界source它只是在等待后续消息。Flink无法预知未来有没有数据所以Kafka source永远不会通知下游“数据结束了”。反过来读取MySQL的一张表如果只做一次全量快照扫描那它就是有界流。但如果用Flink CDC监听这张表的binlog变化那同样是Kafka-like的持续流属于无界流。全量快照是“当前这一秒的表内容”binlog是“从某一秒开始发生的所有变更”一个是有限集合一个是无限事件流。这个区分在后续选型时极其重要。2. 流批一体的底层逻辑一套DAG如何同时跑批和流2.1 从DataStream API到StreamGraph再到JobGraph统一抽象到底统一在哪Flink能实现流批一体不是因为它同时维护了两套引擎而是因为所有任务都先被构建成同一套逻辑图再根据数据源的有界性和运行时配置决定如何执行这张图。具体链路是用户代码 (DataStream API / Table API / SQL) ↓ StreamGraph逻辑层面的算子拓扑只关心“谁转换谁”不关心运行细节 ↓ JobGraph在逻辑图基础上生成的可提交任务单元包含算子链、共享槽等 ↓ ExecutionGraph并行化后的物理执行图每个算子拆成多个并行子任务 ↓ TaskManager上真正执行的Task在这条链路里StreamGraph阶段完全不分批和流。你在SQL里写的SELECT、JOIN、GROUP BY、WINDOW生成的都是一样的算子节点和边。真正区别从JobGraph往下的调度和shuffle部分才体现出来。所以所谓批流一体指的是“同一套API、同一个逻辑图、同一个作业定义”而不是“批处理和流处理的运行方式完全一样”。Flink 1.12开始官方引入RuntimeExecutionMode概念目的就是让你用一套代码声明作业然后按模式自动优化执行。2.2 BATCH模式下Flink到底改了哪些关键部件在STREAMING模式下Flink的调度器会把所有任务一次性地调度起来因为无界流必须全程保活每条数据来都要立刻传递到下游。这也是为什么流模式对资源和稳定性要求极高——所有算子的生命周期是重叠的任何一个task阻塞都会引发背压传导。在BATCH模式下因为输入有界Flink可以把整张图拆成多个调度阶段SchedulingStage按照数据的流动顺序分阶段执行。上游阶段处理完并落盘下游阶段再启动。这种“阶段式调度”有几个直接好处每个阶段的坏节点只需要重启本阶段不需要重启整个作业。shuffle边可以用Blocking模式上游全部算完再一次性传给下游没有背压问题。同一个算子在不同阶段可以利用更优的数据结构比如排序聚合代替哈希聚合。此外BATCH模式默认会关闭很多流处理特有的机制不需要连续生成Checkpoint不需要维持长连接不需要为每个并行子任务维持常驻资源。这也解释了为什么同一个数据量用BATCH模式跑往往比STREAMING模式跑快不少——它根本不是同一个引擎调优的产物而是针对有界数据重新设计了执行计划。2.3 窗口在有界流上并不是难点真正的难点是时间语义很多人一听到流处理就想到窗口然后误以为离线批处理没有窗口概念。其实Flink里窗口API完全共享。在有界流上开一个15分钟的滚动窗口所有数据一次性读入窗口边界可以精确计算连Watermark都不需要人为处理——读到末尾会自动触发所有窗口计算。但在无界流上窗口计算的核心变成了“怎么确定数据截止时间”。数据来了但不能确定后面还有没有属于这个窗口的数据所以要靠Watermark。窗口本身只是个状态容器难点全在“什么时候关闭窗口”以及“迟到的数据怎么处理”。这一点在下一章展开但先记住一个结论批流一体在窗口层面并没有本质冲突批模式只是把“等待窗口关闭”的时间压缩到无限接近于零。3. 时间语义、状态和容错批流一体真正的硬骨头3.1 Watermark和Event Time无界流里的“全局钟表”时间语义有三类Processing Time当前机器时间、Event Time数据里携带的事件时间、Ingestion Time进入Flink的时间。平时排错、调窗口绕不开的是Event Time和Processing Time两套时间体系。Processing Time实现简单延迟低但结果不确定——同一批数据在不同时间重跑窗口归属可能变化。Event Time准确但需要额外机制应对乱序和迟到。Watermark就是用来测量Event Time进度的“水位线”假设一个WatermarkW(t)表示“事件时间小于等于t的数据已经全部到达”那么窗口在收到水位线越过窗口结束时间时就可以安全触发计算。如果数据乱序很严重Watermark需要设置一定的延迟比如WITH DELAY或allowedLateness用延迟换取准确度。在有界流场景Flink会在source结束时自动发送一个Long.MAX_VALUE的Watermark强行把所有未触发的窗口全部推完。所以离线作业根本不需要关心水位线未触发的问题。无界流则没那么幸运任何时刻都可能出现“有几个分区的数据一直没更新”的情况导致Watermark长期不推进。我在实际项目里就遇到过Kafka某个topic有8个分区但只有一个分区持续有数据其他7个分区几分钟才来一条。如果不给source设置withIdleness(Duration.ofSeconds(60))空闲分区超时那么整个Flink任务的Watermark会被空转分区死死卡住窗口永远不触发数据和作业都停在界面上看起来就像“卡死”了。这个问题在批流一体场景下尤其容易被忽视——因为你用离线思维写任务根本不会想到要设置空闲检测。3.2 状态的生命周期从“任务结束即清除”到“无限增长”有界流的状态天然是隐喻的一个批作业聚合一亿条记录无论底层的状态多大作业结束后都可以立即释放。无界流的状态则会伴随作业生命周期无限增长而且因为要不断更新还会产生频繁的增量快照。因此Flink的状态后端设计本质上是为了无界流服务的。RocksDBStateBackend可以把状态刷到磁盘用内存做缓存从而支撑TB级状态HashMapStateBackend把所有状态放堆内存读写快但上限低。在批流一体中如果你在BATCH模式运行一个全局keyBy聚合RocksDB的序列化开销反而可能拖慢速度——这时用HashMap或者干脆用BATCH模式下的排序聚合更合适。状态TTLTime-to-Live也是无界流才有的需求。批处理几小时跑完根本不需要“清理过期Key”。无界流如果设计不合理状态会像漏水的水桶一样无限增长。比如按用户ID维护登录状态如果不设置TTL一年后光状态就能压垮内存。批流一体最阴险的地方就在这你在批处理上能跑通的逻辑搬到无界流上可能一周后OOM原因不是Flink错了而是状态生命周期没有重新思考。3.3 Exactly-Once语义批和流的“保证成本”完全不同Flink的Exactly-Once语义在批模式下几乎“免费”实现有界输入失败部分可以通过阶段重放来重算sink可以等整个阶段完成后一次性写入不需要两阶段提交。无界流则复杂得多。要做到端到端Exactly-Once需要三件套配合Checkpoint保存每个算子的状态快照。两阶段提交协议TwoPhaseCommitSinkFunction让外部sink在Checkpoint完成时原子提交事务。依赖Kafka、Doris、Iceberg等外部系统对事务性写入的支持。所以你会看到市面上很多号称“批流一体”的框架在“一致性”这一层往往是打折的。批可以直接重算流必须回滚后重放虽然结果是“都精确一次”但机制、代价、运维心智完全不是一个量级。实际选型建议如下场景推荐模式理由每日T1报表读取HDFS/文件BATCH有界输入阶段调度省资源实时大屏监控指标STREAMING数据无穷无尽必须持续计算数据湖从MySQL全量同步BATCH快照有界跑完即止数据湖实时增量捕获CDCSTREAMINGbinlog是无界事件流Kafka消费预处理后写入数仓STREAMING数据流无界无法预知结束读取Hive分区表做小时级计算BATCH每个分区是固定集合上表看起来简单但实际项目中经常有人用STREAMING模式去跑一个有界文件源导致任务永远不结束、资源白占也有人用BATCH模式直接跑Kafka源结果数据持续产出但作业从不触发周期性输出。归根结底是你没有在建模阶段想清楚“我的数据源头会不会结束”这件事。4. 从“看得懂”到“用得好”实际场景中的批流一体落地4.1 数据湖上的批流一体当Iceberg遇上Flink CDC数据湖场景是批流一体最典型的战场尤其是近几年选型Iceberg或Hudi的公司。数据湖本质上是一个可低成本重写、带事务能力的文件系统它天然支持批读写也支持流式upsert。一个比较完善的落地方式是用Flink CDC读取数据库的binlog持续捕获变更写入Iceberg表无界流。每日凌晨用Flink BATCH跑一次全量快照同步修正长时间运行可能产生的漂移有界流。下游用Flink SQL直接查Iceberg表既可以跑批量聚合也可以做持续查询。这里要特别提醒Flink CDC本身也是一个“有界转无界”的混合体。CDC任务启动时往往先做一次历史数据的快照有界然后无缝切换到binlog监听无界。你要明白这个切换点才不会在“全量阶段一直不结束”或者“增量阶段丢数据”时抓瞎。建议全量阶段把并行度调低避免对源库压力过大增量阶段再调高并行度和Checkpoint间隔保证时效性。4.2 Flink SQL中怎么用一套SQL跑批和流Flink SQL天然支持批流一体关键在于运行时模式设置。-- 设置为流模式 SET execution.runtime-mode STREAMING; -- 或设置为批模式 SET execution.runtime-mode BATCH; -- 也可以写在提交命令里 flink run -Dexecution.runtime-modeBATCH -c com.example.MyJob app.jarDataStream API对应写法是StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeExecutionMode.BATCH);需要注意的是不是所有连接器都同时支持两种模式。比如JDBC连接器在BATCH模式下默认做一次性的表扫描非常适合批量读取MySQL或PostgreSQL在STREAMING模式下它变成轮询查询有scan.fetch-size和scan.auto-commit参数影响行为。如果你在流模式下游看到数据库连接数异常飙升多半就是轮询间隔配得太短或者fetch-size太小导致连接迟迟不释放。flink sql client/sql gateway里最常见的一个坑是运行批任务时忘了显式声明BATCH模式而SQL Gateway默认宽表配置可能是STREAMING于是明明读的是一个文件或者一个静态表任务却一直不结束。这个问题排查思路很简单先看execution.runtime-mode再看source的scan.bounded.mode。4.3 常见坑的完整排查链路Watermark不触发、JDBC连接器异常、资源估算失误先说Watermark不触发。症状是窗口结果迟迟不出或者只有一部分窗口出结果。排查顺序建议如下确认source是否设置了Watermark生成器以及生成策略是否正确。确认并行度。如果source的并行度大于Watermark生成的并行度某些分区的Watermark可能无法聚合。确认是否有长期不更新的空闲分区。如果是加上withIdleness或table.exec.source.idle-timeout。确认是否把allowedLateness设得过大导致窗口明明触发了却不输出一直在等迟到数据。确认sink的幂等性或提交模式有时窗口结果已经算出来但sink没有提交成功界面看不到输出。再说JDBC连接器异常。我记得很深刻的一次是线上一个STREAMING模式的任务每5秒轮询一次MySQL跑了一个月后开始频繁报“Too many connections”。原因是连接器每轮询都会新建连接又没有正确复用。后来在SQL中加上scan.fetch-size、scan.auto-commit、以及连接池配置并且把轮询间隔从5秒改成30秒才把连接数压下来。建议凡是JDBC源进入STREAMING模式一定要确认下游是否需要如此高频地查源库宁可加一点延迟也要保护数据库。最后说资源估算。批流一体的资源规划不能一拍脑袋。批任务按数据总量估算1TB输入50个并行度每个并行处理20GB内存按每个slot 4GB估算就行。流任务则要看每秒速率、每条数据大小、状态规模、窗口跨度。同样是日均1亿条数据如果集中在凌晨2点导入批模式30分钟跑完也许只用40个CU如果均匀分布在全天流模式可能需要常驻100个CU。这种差别不做资源模型推演根本发现不了。5. 一个可复现的验证实验同一套SQL切换BATCH和STREAMING观察行为差异5.1 准备一个有界和一个无界数据源实验最好准备两个表一个读CSV文件有界一个读Kafka无界。SQL如下-- 有界源读取本地文件 CREATE TABLE csv_source ( user_id STRING, amount DECIMAL(10, 2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector filesystem, path /tmp/orders.csv, format csv ); -- 无界源读取Kafka CREATE TABLE kafka_source ( user_id STRING, amount DECIMAL(10, 2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic orders, properties.bootstrap.servers localhost:9092, properties.group.id test-group, format csv );再建一张输出表CREATE TABLE result_sink ( window_start TIMESTAMP(3), user_id STRING, total_amount DECIMAL(10, 2) ) WITH ( connector print );聚合查询INSERT INTO result_sink SELECT TUMBLE_START(ts, INTERVAL 1 MINUTE), user_id, SUM(amount) FROM csv_source GROUP BY TUMBLE(ts, INTERVAL 1 MINUTE), user_id;5.2 运行观察调度方式完全不一样先用BATCH模式跑csv_source你会看到作业从Submitted到Running到Finished整个过程很快日志里不会出现持续的Checkpoint所有窗口数据读完后一次性触发计算。在Flink Web UI的JobGraph上你会发现多个阶段依次执行任务不是所有节点同时启动的。再用STREAMING模式跑kafka_source操作方式基本一样但作业会一直处于Running状态日志里会出现周期性的Checkpoint完成信息。如果Kafka中暂时没有数据Watermark不会无限增大窗口也不会有输出而一旦你往Kafka里发几条延迟超过5秒的消息窗口就会立刻按Watermark语义计算并输出结果。5.3 这个实验告诉我们的三件事第一批和流在Flink里不是两套东西而是同一套SQL的两种执行策略。你完全可以用同一份代码只在运行时切换模式跑通两种场景。第二有界无界这件事不是“代码长得像流就是无界”而是“数据源完不完了才算”。无界流必须长期运行所以要和状态、Checkpoint、Watermark做好朋友。第三实际项目中很多“实时任务”实际上用BATCH模式加短周期调度比如每10分钟调度一次就能满足需求成本和稳定性都更好。反之一些“离线任务”因为业务要求分钟级可见性也应该考虑用无界流来处理。判断标准只有一条你的数据是有限的还是无穷无尽的。在我自己团队的实践中有一类任务以前全用STREAMING模式挂着从Kafka读取用户行为日志做1小时的滚动聚合写入MySQL。后来发现这个需求的时效性其实允许10分钟延迟于是直接改成每10分钟启动一个BATCH作业去读Kafka中过去一个小时的落盘数据。表面上没变但资源占用下降了70%任务失败率也大幅降低。这条经验不一定要照搬但它印证了批流一体的价值不是逼你在批和流里二选一而是让你有能力在同一个技术栈里按数据本身的形状和业务时效要求做最合适的执行选择。
返回列表