
在传统大数据架构中批量计算与流式计算长期由两套独立引擎承载批处理依赖 MapReduce/Spark 等离线引擎流处理依赖 Storm/Spark Streaming/Flink 等实时引擎。这种分离导致同一套业务逻辑需要分别用两套 API 实现语义难以对齐运维成本成倍增加。Apache Flink 的关键设计决策之一是用「有界流Bounded Stream」与「无界流Unbounded Stream」这两个概念将批与流统一到同一套 Runtime、状态管理与容错机制之上。理解这两个概念的精确语义以及它们如何决定算子的执行模型是正确使用 Flink、合理选择执行模式的前提。一、有界流与无界流的形式化定义在 Flink 的 DataStream 模型中一条流Stream被定义为数据记录Record的有序序列。根据该序列是否具有确定的终止条件流被划分为两类有界流Bounded Stream元素集合有限存在明确的起点与终点。作业读取完最后一个元素后即可结束。典型来源HDFS 文件、Hive 表、关系型数据库的全量快照、对象存储S3/OSS中的文件。无界流Unbounded Stream元素集合无限只有起点、不存在自然终点数据持续到达。作业必须常驻运行无法等待「全部数据到齐」。典型来源Kafka/Pulsar/Kinesis 消息队列、传感器数据、用户行为日志。二者的判断标准只有一条数据源是否存在自然终点。这一判定与数据的物理形态无关——一个 CSV 文件被读取后同样表现为一条流只是这条流有界一个 Kafka topic 被消费时也是一条流只是无界。这一概念在 Flink 源码层面有明确对应Source接口的getBoundedness()方法返回Boundedness.BOUNDED或Boundedness.CONTINUOUS_UNBOUNDED它是execution.runtime-mode AUTOMATIC模式判断执行方式的依据。有界流与无界流的本质差异直接决定了二者在算子语义上的根本区别这是下一节的核心。二、有界/无界如何决定算子语义与执行模型有界流与无界流最关键的区别在于它们允许使用的算子类型不同。Flink 的算子Operator按数据消费方式可分为两类流水线式算子Pipeline Operator边接收边处理数据到达即向下游输出无需等待全部输入。map、filter、flatMap等大多数转换属于此类。阻塞式算子Blocking Operator必须消费完全部输入才能产生输出。典型如排序Sort、非窗口的全局聚合Global Aggregation、以及Hash Join 的 Build 端构建。有界流可以使用阻塞式算子。因为数据总量有限引擎可以完整读入全部数据后再执行排序、全局聚合、广播 Join 等操作。这正是批处理Batch的语义基础。无界流只能使用流水线式算子。因为数据永不终止任何「等待全部数据」的算子都无法完成。无界流的计算必须依赖以下三类机制来将无限数据转化为有限的可计算结果窗口Window将无限数据流按时间或数量切分为有限的块对每个块独立聚合。按时间划分时常见类型包括滚动窗口Tumbling、滑动窗口Sliding与会话窗口Session。Watermark声明事件时间的推进进度用于判断某个事件时间窗口是否可以触发、迟到的数据应如何处理。状态State保存算子计算的中间结果支持跨窗口、跨记录的增量累积。Flink 将状态划分为 Keyed State 与 Operator State前者与 Key 绑定、随数据分布在各并行子任务后者与算子实例绑定。在时间语义上Flink 区分三种时间事件时间Event Time数据产生的时间、处理时间Processing Time算子处理该数据的时间与摄取时间Ingestion Time数据进入 Flink 的时间。无界流的窗口聚合必须基于事件时间并配合 Watermark 才能获得确定、可复现的结果处理时间虽实现简单但结果会随处理速度波动且无法应对乱序与迟到数据。由此引出两个最常见的误区误区一在无界流上调用collect()或count()拉取全量数据。这类操作会把数据拉回 Driver 端内存在有界流上可行在无界流上则会导致数据永远收集不完、Driver 端堆内存溢出OOM。误区二在有界流上配置了事件时间 Watermark窗口永不触发。有界流读完后 Watermark 不再推进事件时间窗口因等不到「Watermark 越过窗口结束时间」而无法触发表现为作业执行完毕却不产出结果。三、Flink 批流一体的实现与演进早期 Flink1.12 之前提供两套 APIDataSet API用于批处理DataStream API用于流处理二者底层执行模型不完全一致。这带来的直接问题是同一业务需分别实现两套代码语义对齐困难。Flink 社区随后统一了认识批处理不过是有界流的一种特例。自 Flink 1.12 起批处理被统一到 DataStream API 之上DataSet API标记为废弃并在 1.18 版本彻底移除。自此无论批还是流都运行于同一套 RuntimeJobManager 负责作业调度、资源分配与 Checkpoint 协调TaskManager 负责执行数据流算子、维护状态与网络 Shuffle。状态管理、容错机制Checkpoint/Savepoint与调度逻辑对批流完全一致。批与流的切换通过一个配置项完成而非更换 API# 批执行模式有界流走阻塞式算子支持排序/全局聚合bin/flink run -Dexecution.runtime-modeBATCH ./app.jar# 流执行模式默认按事件流增量处理bin/flink run -Dexecution.runtime-modeSTREAMING ./app.jar# 自动模式有界源自动走批无界源自动走流bin/flink run -Dexecution.runtime-modeAUTOMATIC ./app.jarBATCH模式并非仅是「跑完就退出」的语义差别它在执行层引入了实质性优化启用Sort-based Shuffle以排序方式组织数据交换替代流模式的 Pipeline Shuffle、允许阻塞式算子排序、全局聚合、Hash Join、并采用更粗粒度的调度。这些优化在数据规模大、需要全局排序或全量 Join 的场景下能显著降低网络与内存开销。四、代码同一逻辑的批流两种执行方式下面以「统计订单金额」为例分别演示有界流与无界流的实现。依赖 Flink 1.18代码可直接编译运行。4.1 有界流读取文件执行结束后退出importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;publicclassBoundedStreamDemo{publicstaticvoidmain(String[]args)throwsException{// 默认 STREAMING 模式因数据源有界处理完最后一条记录后作业自然结束finalStreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// readTextFile 读取一个有界流文件有终点读完即止DataStreamStringlinesenv.readTextFile(hdfs:///data/orders.log);// 每行格式: order_id,amount,event_time,category取第 2 列金额DataStreamDoubleamountslines.map(line-Double.parseDouble(line.split(,)[1]));// 有界流可用全局聚合数据全量可读sum 在读取完毕时输出最终结果amounts.keyBy(v-total).sum(0).print();env.execute(bounded-stream-demo);}}要点虽然环境为 STREAMING 模式但因数据源有界sum算子能在读取完毕时给出确定结果并退出。若将其改为BATCH模式Flink 会进一步启用阻塞式聚合与 Sort-based Shuffle 优化。4.2 无界流读取 Kafka基于事件时间窗口聚合importorg.apache.flink.api.common.eventtime.WatermarkStrategy;importorg.apache.flink.api.common.serialization.SimpleStringSchema;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;importorg.apache.flink.streaming.api.windowing.time.Time;importorg.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;importjava.time.Duration;importjava.util.Properties;publicclassUnboundedStreamDemo{// 订单事件 POJOFlink 的 sum(amount) 要求字段可被反射访问publicstaticclassOrder{publiclongorderId;publicdoubleamount;publiclongeventTime;// 事件发生时间毫秒时间戳publicStringcategory;publicOrder(){}publicOrder(longorderId,doubleamount,longeventTime,Stringcategory){this.orderIdorderId;this.amountamount;this.eventTimeeventTime;this.categorycategory;}publicstaticOrderparse(Stringline){String[]pline.split(,);returnnewOrder(Long.parseLong(p[0]),Double.parseDouble(p[1]),Long.parseLong(p[2]),p[3]);}}publicstaticvoidmain(String[]args)throwsException{finalStreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();PropertieskafkaPropsnewProperties();kafkaProps.setProperty(bootstrap.servers,localhost:9092);kafkaProps.setProperty(group.id,flink-order-group);// 无界源Kafka topic 持续产生数据无自然终点DataStreamStringrawenv.addSource(newFlinkKafkaConsumer(orders,newSimpleStringSchema(),kafkaProps));DataStreamDoubletotalsraw.map(Order::parse)// 声明事件时间 允许 5 秒有界乱序的 Watermark.assignTimestampsAndWatermarks(WatermarkStrategy.OrderforBoundedOutOfOrderness(Duration.ofSeconds(5)).withTimestampAssigner((order,ts)-order.eventTime))// 按类别分组使聚合分布到各并行子任务.keyBy(order-order.category)// 事件时间滚动窗口每 1 分钟聚合一次.window(TumblingEventTimeWindows.of(Time.minutes(1))).sum(amount);totals.print();// 常驻运行不主动退出env.execute(unbounded-stream-demo);}}要点无界流不存在「最终结果」只有「截至某个事件时间点的结果」该时间点由 Watermark 决定。forBoundedOutOfOrderness(Duration.ofSeconds(5))声明允许 5 秒乱序窗口在 Watermark 越过其结束时间时触发。4.3 同一段 SQL批流通吃Table API / SQL通过 Table API / SQL批流切换进一步简化——逻辑只编写一次连接的表决定了执行方式-- 有界表连接文件批执行CREATETABLEorders_bounded(order_idBIGINT,amountDOUBLE,tsTIMESTAMP(3))WITH(connectorfilesystem,pathhdfs:///data/orders,formatcsv);-- 无界表连接 Kafka流执行声明 WatermarkCREATETABLEorders_unbounded(order_idBIGINT,amountDOUBLE,tsTIMESTAMP(3),WATERMARKFORtsASts-INTERVAL5SECOND)WITH(connectorkafka,topicorders,properties.bootstrap.serverslocalhost:9092,formatjson);-- 同一句聚合运行于有界表即为批运行于无界表即为流SELECTwindow_start,SUM(amount)AStotalFROMTABLE(TUMBLE(TABLEorders_unbounded,DESCRIPTOR(ts),INTERVAL1MINUTE))GROUPBYwindow_start;五、执行模式切换的语义差异与误区execution.runtime-mode的三个取值具有明确的语义边界模式适用数据源算子语义典型表现STREAMING无界/有界均可流水线式增量处理作业常驻或读完后退出BATCH仅限有界源阻塞式支持排序/全局聚合读完后输出结果并退出AUTOMATIC二者均可依据Source.getBoundedness()自动选择有界走批、无界走流基于此有以下四个需要规避的误区误区一在无界流上调用collect()拉取全量数据。无界流数据永无穷尽collect()将数据持续拉回 Driver 端内存必然导致 OOM。排查数据应改用受限的print()、或写入 Sink 后从下游存储读取。误区二在有界流上配置事件时间 Watermark窗口不触发。数据读完后 Watermark 停止推进事件时间窗口无法满足触发条件。有界流若需窗口聚合应改用处理时间或切换为BATCH模式由引擎按有界语义优化。误区三runtime-mode设为BATCH但 Source 为无界源。BATCH模式要求数据源有界可读完否则作业在提交阶段即报错。反之若 Source 为有界源但设为STREAMING语义仍正确只是无法享受批执行优化。误区四AUTOMATIC模式遇到未正确实现getBoundedness()的自定义/第三方 Connector静默退化为流。若 Connector 未正确上报有界性AUTOMATIC会将其按无界流处理表现为「作业跑完不退出」。因此凡能明确判定的场景应显式指定STREAMING或BATCH而非依赖AUTOMATIC的自动推断。六、执行模式选型建议针对常见业务场景给出如下判断离线报表、历史数据回填、全量初始化→ 有界源 BATCH模式。数据总量有限引擎可启用排序、全局聚合、Sort-based Shuffle 等优化兼顾性能与资源开销。实时告警、实时大屏、实时数仓→ 无界源 STREAMING模式配合事件时间窗口、Watermark 与状态完成增量聚合。同一逻辑先离线回填历史、再流处理增量→ 这是 Flink 批流一体的核心价值。用 Table API / SQL 编写一次逻辑历史数据连接有界表以BATCH执行增量数据连接无界表以STREAMING执行两段共享同一份代码与语义。有界流与无界流的本质区别不在于数据的物理形态而在于数据源是否存在自然终点这一区别进一步决定了算子可采用流水线式还是阻塞式执行模型从而划分出批与流的语义边界。掌握了这一底层逻辑execution.runtime-mode的选型便有了明确依据。同一逻辑先离线回填历史、再流处理增量** → 这是 Flink 批流一体的核心价值。用 Table API / SQL 编写一次逻辑历史数据连接有界表以BATCH执行增量数据连接无界表以STREAMING执行两段共享同一份代码与语义。