在传统大数据架构中,批量计算与流式计算长期由两套独立引擎承载:批处理依赖 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 批流一体的实现与演进
早期 Flink(1.12 之前)提供两套 API:DataSet API用于批处理,DataStream API用于流处理,二者底层执行模型不完全一致。这带来的直接问题是:同一业务需分别实现两套代码,语义对齐困难。
Flink 社区随后统一了认识:批处理不过是有界流的一种特例。自 Flink 1.12 起,批处理被统一到 DataStream API 之上,DataSet API标记为废弃,并在 1.18 版本彻底移除。自此,无论批还是流,都运行于同一套 Runtime:JobManager 负责作业调度、资源分配与 Checkpoint 协调,TaskManager 负责执行数据流算子、维护状态与网络 Shuffle。状态管理、容错机制(Checkpoint/Savepoint)与调度逻辑对批流完全一致。
批与流的切换通过一个配置项完成,而非更换 API:
# 批执行模式:有界流走阻塞式算子,支持排序/全局聚合bin/flink run -Dexecution.runtime-mode=BATCH ./app.jar# 流执行模式(默认):按事件流增量处理bin/flink run -Dexecution.runtime-mode=STREAMING ./app.jar# 自动模式:有界源自动走批,无界源自动走流bin/flink run -Dexecution.runtime-mode=AUTOMATIC ./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 模式;因数据源有界,处理完最后一条记录后作业自然结束finalStreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();// readTextFile 读取一个有界流:文件有终点,读完即止DataStream<String>lines=env.readTextFile("hdfs:///data/orders.log");// 每行格式: order_id,amount,event_time,category,取第 2 列金额DataStream<Double>amounts=lines.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{// 订单事件 POJO;Flink 的 sum("amount") 要求字段可被反射访问publicstaticclassOrder{publiclongorderId;publicdoubleamount;publiclongeventTime;// 事件发生时间(毫秒时间戳)publicStringcategory;publicOrder(){}publicOrder(longorderId,doubleamount,longeventTime,Stringcategory){this.orderId=orderId;this.amount=amount;this.eventTime=eventTime;this.category=category;}publicstaticOrderparse(Stringline){String[]p=line.split(",");returnnewOrder(Long.parseLong(p[0]),Double.parseDouble(p[1]),Long.parseLong(p[2]),p[3]);}}publicstaticvoidmain(String[]args)throwsException{finalStreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();PropertieskafkaProps=newProperties();kafkaProps.setProperty("bootstrap.servers","localhost:9092");kafkaProps.setProperty("group.id","flink-order-group");// 无界源:Kafka topic 持续产生数据,无自然终点DataStream<String>raw=env.addSource(newFlinkKafkaConsumer<>("orders",newSimpleStringSchema(),kafkaProps));DataStream<Double>totals=raw.map(Order::parse)// 声明事件时间 + 允许 5 秒有界乱序的 Watermark.assignTimestampsAndWatermarks(WatermarkStrategy.<Order>forBoundedOutOfOrderness(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('connector'='filesystem','path'='hdfs:///data/orders','format'='csv');-- 无界表:连接 Kafka,流执行,声明 WatermarkCREATETABLEorders_unbounded(order_idBIGINT,amountDOUBLE,tsTIMESTAMP(3),WATERMARKFORtsASts-INTERVAL'5'SECOND)WITH('connector'='kafka','topic'='orders','properties.bootstrap.servers'='localhost:9092','format'='json');-- 同一句聚合,运行于有界表即为批,运行于无界表即为流SELECTwindow_start,SUM(amount)AStotalFROMTABLE(TUMBLE(TABLEorders_unbounded,DESCRIPTOR(ts),INTERVAL'1'MINUTE))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执行,两段共享同一份代码与语义。