Article

数据计算 Flink

更新于:2026-07-13

1.1 什么是 Flink:流处理与批处理统一

概念名称说明注意事项
流处理(Streaming Processing)Flink 的核心是流处理,将所有数据视为无限流(unbounded stream),支持低延迟、高吞吐的实时计算。流处理适用于实时监控、日志分析、风控等场景。
批处理(Batch Processing)批处理被视为流处理的特例——有界流(bounded stream)。Flink 使用同一引擎处理批和流,无需切换框架。统一引擎减少运维复杂度,但批处理性能需合理配置并行度和内存。
事件驱动应用Flink 可构建事件驱动系统,如基于规则的告警、状态机转换等。需结合状态管理和时间语义实现复杂逻辑。
数据管道(Data Pipeline)支持数据在系统间高效流转,如 ETL、数据清洗、格式转换等。可通过 Savepoint 实现应用升级与回滚。
流批一体架构Flink 提供统一 API(DataStream API 和 Table/SQL)处理流与批,底层运行时一致。建议新项目优先使用 DataStream API + 事件时间,避免 DataSet API(已弃用)。
组件名称说明注意事项
JobManager负责调度任务、协调检查点(Checkpoint)、管理元数据。集群中可有多个,但只有一个为 Leader。高可用模式下需配合 ZooKeeper 或 Kubernetes 实现故障转移。
TaskManager实际执行任务的工作节点,管理内存、网络、任务线程。多个 TaskManager 形成资源池。每个 Task Slot 执行一个或多个算子子任务,Slot 数量影响并行度上限。
Client提交作业的客户端,负责解析程序、生成 JobGraph 并提交给 JobManager。Client 可运行在本地或远程,提交后可断开(detached 模式)。
Task SlotTaskManager 中的资源单位,每个 Slot 独立执行一个任务链(task chain)。Slot 数量由内存决定,不隔离 CPU,建议根据资源需求合理配置。
JobGraph / ExecutionGraphClient 将程序转换为 JobGraph,JobManager 转为 ExecutionGraph 调度执行。JobGraph 是逻辑图,ExecutionGraph 是运行时视图。

1.3 数据流模型:DataStream 与 DataSet(旧)

模型名称说明注意事项
DataStream API核心 API,用于处理无界和有界数据流,支持事件时间、窗口、状态等高级特性。推荐用于所有新项目,支持流批统一处理。
DataSet API旧版批处理 API,仅支持有界数据集,功能已被 DataStream 和 Table API 取代。自 Flink 1.12 起已标记为过时(deprecated),不建议新项目使用。
DataStream表示一个数据流,可以是无界(流)或有界(批),通过 env.addSource() 创建。所有转换操作(如 map、filter)返回新的 DataStream。
Bounded vs UnboundedFlink 将有界流视为无界流的特例,统一调度和执行。有界流作业执行完毕后自动终止,无界流持续运行。
Transformation数据转换操作,如 map、keyBy、window 等,构建 DAG 图。转换是懒加载的,只有调用 execute() 才触发执行。
阶段说明注意事项
1. 编写程序用户使用 DataStream API 或 Table API 编写逻辑。程序需定义执行环境(StreamExecutionEnvironment)。
2. 创建 StreamGraph客户端根据 API 调用生成逻辑图(StreamGraph)。图中节点为算子,边为数据流。
3. 生成 JobGraph客户端优化 StreamGraph(如任务链合并),生成 JobGraph 提交至 JobManager。JobGraph 是调度的基本单位。
4. 生成 ExecutionGraphJobManager 将 JobGraph 转为 ExecutionGraph,分配任务到 TaskManager。ExecutionGraph 包含任务的并行实例(subtask)。
5. 调度与执行TaskManager 启动任务,开始消费数据并执行计算。支持 Checkpoint 机制保障容错。
6. 结果输出数据通过 Sink 输出到外部系统(如 Kafka、文件、数据库)。Sink 可配置异步或同步写入。
部署模式说明注意事项
Local 模式单机运行,JobManager 和 TaskManager 在同一 JVM 中,用于开发测试。无需配置,直接运行 main 方法即可。
Standalone 模式独立集群,手动启动 JobManager 和 TaskManager,适合固定资源环境。需手动管理高可用(HA)和资源伸缩。
YARN 模式基于 Hadoop YARN 资源调度,支持 per-job 和 session 模式。需 Hadoop 环境,适合已有 YARN 的企业。
Kubernetes 模式在 K8s 上部署 Flink 集群,支持动态伸缩和高可用。推荐生产环境使用,需熟悉 K8s 配置。
Session Cluster多个作业共享一个集群,资源利用率高,但作业间可能干扰。不推荐生产使用,隔离性差。
Per-Job Cluster每个作业独占集群,资源隔离好,启动开销大。推荐生产环境使用,尤其 YARN/K8s 模式。
Application Mode作业代码打包提交,集群由 Flink 管理,代码在集群内启动。与 Per-Job 类似,但更符合云原生理念。

第二章:开发环境搭建与项目初始化

2.1 搭建本地开发环境(Java/Scala)

组件说明注意事项
JDK 8 或 11Flink 基于 JVM,需安装 Java 8 或 11(推荐 11)。不支持 Java 17+(截至 Flink 1.17),建议使用 OpenJDK。
Maven 3.0+用于依赖管理和项目构建。需配置镜像(如阿里云)加速下载。
Flink 版本选择推荐使用最新稳定版(如 1.17 或 1.18)。避免使用过旧版本,确保安全与功能支持。
Scala 版本(如用 Scala)Flink Scala API 需匹配 Scala 版本(如 2.12 或 2.13)。Java 项目可忽略 Scala 依赖。
环境变量建议设置 JAVA_HOME,Maven 自动识别。Windows 用户注意路径分隔符。
配置项语法用途代码示例注意事项
groupId / artifactId<groupId>com.example</groupId>
<artifactId>flink-demo</artifactId>
定义项目坐标见完整 pom 示例建议使用反向域名命名
Flink 版本属性<properties>
<flink.version>1.17.0</flink.version>
</properties>
统一管理 Flink 版本同上避免版本冲突
Flink 依赖引入<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
引入核心 API同上必须引入 streaming 模块
Scala 依赖(可选)<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-scala_2.12</artifactId>
<version>${flink.version}</version>
</dependency>
使用 Scala API同上注意 Scala 版本匹配
打包插件<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.4.1</version>
<executions>...
<configuration>
<shadedArtifactAttached>true</shadedArtifactAttached>
</configuration>
</plugin>
打成 fat jar,包含所有依赖同上提交到集群必须使用 fat jar
程序类型说明代码示例注意事项
批式 WordCount使用 ExecutionEnvironment 处理有界数据final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
DataSet<String> text = env.fromElements(
"hello world", "hello flink", "flink is awesome"
);
DataSet<Tuple2<String, Integer>> counts = text
.flatMap((String line, Collector<String> out) -> {
Arrays.stream(line.split(" ")).forEach(out::collect);
})
.map(word -> Tuple2.of(word, 1))
.returns(Types.TUPLE(String.class, Integer.class))
.groupBy(0)
.sum(1);
counts.print();
DataSet API 已弃用,仅用于演示
流式 WordCount使用 StreamExecutionEnvironment 处理流数据final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = env.fromElements(
"hello world", "hello flink", "flink is awesome"
);
DataStream<Tuple2<String, Integer>> counts = text
.flatMap((String line, Collector<String> out) -> {
Arrays.stream(line.split(" ")).forEach(out::collect);
})
.returns(Types.STRING)
.map(word -> Tuple2.of(word, 1))
.returns(Types.TUPLE(String.class, Integer.class))
.keyBy(t -> t.f0)
.sum(1);
counts.print();
env.execute("WordCount");
必须调用 execute() 触发执行
execute() 方法触发程序执行,阻塞直到完成或出错env.execute("Job Name");参数为作业名称,显示在 Web UI

2.4 IDE 配置与调试技巧(IntelliJ IDEA)

配置项说明注意事项
创建 Maven 项目使用 maven-archetype-quickstart 或手动配置 pom.xml确保打包类型为 jar
导入 Flink 依赖Maven 自动下载依赖,IDEA 提示”Import Changes”时确认若依赖未加载,点击刷新按钮
主类运行配置创建 Run Configuration,选择主类(含 main 方法)可直接运行流/批程序,无需集群
断点调试在 map、flatMap 等函数中设置断点,查看数据流支持 Lambda 调试,注意变量捕获
日志查看控制台输出日志,或配置 log4j.properties调试时建议设置日志级别为 INFO 或 DEBUG
本地并行度设置默认并行度为 CPU 核数,可通过 env.setParallelism(1) 调整调试时建议设为 1,便于跟踪
避免依赖冲突排除传递依赖(如旧版本 Jackson)使用 mvn dependency:tree 检查

第三章:DataStream API 基础操作

3.1 创建执行环境(ExecutionEnvironment / StreamExecutionEnvironment)

方法名称语法用途代码示例注意事项
getExecutionEnvironment(批)ExecutionEnvironment.getExecutionEnvironment()创建批处理执行环境,自动判断运行模式(本地或集群)final ExecutionEnvironment batchEnv = ExecutionEnvironment.getExecutionEnvironment();适用于 DataSet API(已弃用)
getExecutionEnvironment(流)StreamExecutionEnvironment.getExecutionEnvironment()创建流处理执行环境,推荐用于所有流式应用final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();自动识别运行环境,开发调试最常用
createLocalEnvironmentStreamExecutionEnvironment.createLocalEnvironment()创建本地执行环境,指定并行度final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(2);用于本地测试,限制并行度为 2
createRemoteEnvironmentStreamExecutionEnvironment.createRemoteEnvironment(String host, int port, String... jarFiles)连接远程 Flink 集群提交作业final StreamExecutionEnvironment env = StreamExecutionEnvironment.createRemoteEnvironment(
"jobmanager", 8081, "path/to/job.jar"
);
需指定 JobManager 地址和 JAR 路径
setParallelismenv.setParallelism(int parallelism)设置整个程序的默认并行度env.setParallelism(4);可被算子级并行度覆盖
disableOperatorChainingenv.disableOperatorChaining()禁用任务链(chaining),每个算子独立调度env.disableOperatorChaining();调试时便于观察子任务
configureenv.getConfig().set...配置运行时参数,如类型序列化、全局作业参数env.getConfig().setAutoWatermarkInterval(200);常用于设置水印生成间隔

3.2 数据源(Source):从集合、文件、Socket、Kafka 读取数据

Source 类型语法用途代码示例注意事项
fromElementsenv.fromElements(T... elements)从 Java 元素创建数据流,用于测试DataStream<String> stream = env.fromElements("a", "b", "c");数据量小,适合本地调试
fromCollectionenv.fromCollection(Collection<T> collection)从集合创建流List<String> data = Arrays.asList("x", "y");
DataStream<String> stream = env.fromCollection(data);
支持任意 Collection 类型
readTextFileenv.readTextFile(String path)读取文本文件,每行作为一个元素DataStream<String> lines = env.readTextFile("/tmp/input.txt");路径支持本地或 HDFS,文件必须有界
socketTextStreamenv.socketTextStream(String hostname, int port)从 Socket 读取文本流,按行分割DataStream<String> stream = env.socketTextStream("localhost", 9999);仅用于测试,生产环境不推荐
addSourceenv.addSource(SourceFunction<T> function)添加自定义数据源env.addSource(new FlinkKafkaConsumer<>(
"topic", new SimpleStringSchema(), props));
通用接口,支持 Kafka、MySQL Binlog 等
FlinkKafkaConsumernew FlinkKafkaConsumer<>(topic, deserializationSchema, props)从 Kafka 读取数据Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "flink-group");
DataStream<String> stream = env.addSource(
new FlinkKafkaConsumer<>("topic", Schema.STRING(), props)
);
需引入 flink-connector-kafka 依赖

3.3 基本转换操作:map、flatMap、filter、keyBy

转换操作语法用途代码示例注意事项
mapstream.map(MapFunction<T, R>)将每个元素转换为另一个元素DataStream<Integer> doubled = stream.map(x -> x * 2);返回一对一映射结果
flatMapstream.flatMap(FlatMapFunction<T, R>)将每个元素转换为零或多个元素stream.flatMap((String line, Collector<String> out) -> {
for (String word : line.split(" ")) {
out.collect(word);
}
})
.returns(Types.STRING);
常用于分词、展开嵌套结构
filterstream.filter(FilterFunction<T>)过滤不符合条件的元素DataStream<String> filtered = stream.filter(s -> !s.isEmpty());返回布尔值,true 保留,false 丢弃
keyBystream.keyBy(KeySelector<T, K>)按键对数据进行分区,用于后续聚合KeyedStream<Tuple2<String, Integer>, String> keyed = stream.keyBy(t -> t.f0);必须作用于 DataStream,返回 KeyedStream
keyBy 字段位置stream.keyBy(0)按元组第 0 个字段分组(旧语法)stream.keyBy(0).sum(1);仅适用于 Tuple 类型,建议使用 Lambda
returnsflatMap(...).returns(TypeInformation)显式声明泛型类型信息.returns(Types.STRING)使用 Lambda 时防止类型擦除,必须调用

3.4 输出操作(Sink):打印、写入文件、Kafka、自定义 Sink

Sink 类型语法用途代码示例注意事项
print / printToErrstream.print() / stream.printToErr()将数据打印到标准输出或错误流stream.print();
stream.print("prefix: ");
输出带并行子任务 ID,如 1> hello
writeAsTextstream.writeAsText(String path)将数据写入文本文件stream.writeAsText("/tmp/output");仅用于有界流,生产环境不推荐
addSinkstream.addSink(SinkFunction<T> function)添加自定义 Sinkstream.addSink(new FlinkKafkaProducer<>(
"topic", new SimpleStringSchema(), props));
通用接口,支持多种外部系统
FlinkKafkaProducernew FlinkKafkaProducer<>(topic, serializationSchema, props)写入 Kafka 主题Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
stream.addSink(new FlinkKafkaProducer<>(
"result-topic", new SimpleStringSchema(), props));
需引入 Kafka Connector 依赖
自定义 Sinkclass MySink implements SinkFunction<T> { ... }实现自定义输出逻辑public class MySink implements SinkFunction<String> {
@Override
public void invoke(String value, Context ctx) {
System.out.println("Custom: " + value);
}
}
stream.addSink(new MySink());
可结合数据库连接、HTTP 调用等

3.5 执行程序与设置并行度

方法名称语法用途代码示例注意事项
executeenv.execute(String jobName)触发程序执行,启动作业env.execute("My Flink Job");必须调用,否则程序不运行
setParallelismenv.setParallelism(int p)设置全局并行度env.setParallelism(4);可被算子级并行度覆盖
setParallelism(算子级)stream.map(...).setParallelism(2)为特定算子设置并行度stream.map(x -> x + 1).setParallelism(2);优先级高于全局设置
getMaxParallelismenv.setMaxParallelism(int max)设置最大并行度,影响状态后端分片env.setMaxParallelism(128);影响 Checkpoint 和状态重平衡
disableOperatorChainingenv.disableOperatorChaining()禁用任务链合并env.disableOperatorChaining();调试时便于观察每个算子
精细控制任务链stream.map(...).disableChaining() / .startNewChain()精细控制任务链stream.map(func).disableChaining();
stream.filter(f).startNewChain();
disableChaining() 断开链,startNewChain() 开始新链

第四章:时间语义与水位线(Watermark)

4.1 事件时间(Event Time)、处理时间(Processing Time)、摄入时间(Ingestion Time)

时间类型说明用途注意事项
事件时间(Event Time)事件在设备上产生的时间,嵌入数据中支持基于真实时间的窗口计算,处理乱序数据需提取时间戳并生成 Watermark
处理时间(Processing Time)数据在 Flink 算子中被处理的系统时间简单高效,延迟最低无法处理乱序或延迟数据,结果不可重现
摄入时间(Ingestion Time)数据进入 Flink Source 的时间折中方案,Source 自动打时间戳不依赖事件自带时间,但无法处理 Source 前的乱序
选择建议env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)设置时间语义Flink 1.12+ 推荐通过 WatermarkStrategy 设置

4.2 为什么需要 Watermark

概念说明注意事项
Watermark(水位线)表示事件时间的进度,是一个带有时间戳的特殊记录用于触发窗口计算,允许一定延迟
乱序数据处理现实中事件可能因网络延迟等原因乱序到达Watermark 告诉系统”早于该时间的事件已基本到达”
延迟容忍Watermark 可设置延迟时间(如 5 秒),允许迟到数据超过延迟的数据默认丢弃,可通过侧输出捕获
触发窗口当 Watermark ≥ 窗口结束时间,触发窗口计算是事件时间窗口的核心机制
单调递增Watermark 必须单调不减,否则无效Source 端生成,沿数据流传播

4.3 Watermark 生成策略:周期性与标点式

策略类型说明适用场景注意事项
周期性 Watermark(Periodic)每隔固定时间(如 200ms)生成 Watermark大多数场景推荐使用通过 assignTimestampsAndWatermarks(WatermarkStrategy) 设置
标点式 Watermark(Punctuated)每当遇到特定事件(如”END”标记)时生成数据流中有明确分界符实现 SourceFunction 时调用 ctx.emitWatermark()
BoundedOutOfOrderness允许最大乱序时间(如 5 秒)乱序程度可预估的场景WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(5))
Ascending Timestamps时间戳严格递增,Watermark = 最大时间戳 - 1几乎无乱序的数据WatermarkStrategy.<String>forMonotonousTimestamps()
No Watermarks不生成 Watermark,仅用于测试调试或处理时间场景WatermarkStrategy.noWatermarks()

4.4 设置事件时间与 Watermark 的 API 使用

方法/类语法用途代码示例注意事项
WatermarkStrategyWatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5))创建 Watermark 策略stream.assignTimestampsAndWatermarks(
WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> extractTime(event))
);
Flink 1.11+ 推荐方式
withTimestampAssigner.withTimestampAssigner((event, timestamp) -> ...)提取事件时间戳.withTimestampAssigner((String event, long timestamp) -> {
return parseTimestamp(event);
})
Lambda 必须返回 long 类型时间戳(毫秒)
assignTimestampsAndWatermarksstream.assignTimestampsAndWatermarks(strategy)应用时间戳和 Watermark 策略见上例必须在 keyBy 和 window 之前调用
TimestampAssigner自定义时间戳提取器从复杂对象中提取时间new SerializableTimestampAssigner<String>() {
@Override
public long extractTimestamp(String element, long recordTimestamp) {
return parseTime(element);
}
}
旧 API,仍兼容
WatermarkGenerator自定义 Watermark 生成逻辑实现复杂乱序处理策略实现 WatermarkGenerator 接口高级用法,需谨慎

第五章:窗口(Window)计算

5.1 窗口类型:滚动窗口、滑动窗口、会话窗口、全局窗口

窗口类型说明适用场景注意事项
滚动窗口(Tumbling Window)固定长度、无重叠的窗口,如每5秒统计一次周期性汇总,如每分钟 PV 统计时间对齐,window(TumblingEventTimeWindows.of(Time.seconds(5)))
滑动窗口(Sliding Window)固定长度、可重叠的窗口,滑动步长小于窗口大小近实时监控,如每1秒输出过去5秒的平均值资源消耗高,因重复计算
会话窗口(Session Window)基于活动间隙(gap)划分,非活跃期结束当前窗口用户行为分析,如会话时长统计支持动态 gap,自动合并邻近会话
全局窗口(Global Window)所有数据分配到一个窗口,需自定义触发器才能输出特殊聚合场景,如全量去重默认不触发,必须配合 Trigger 使用

5.2 窗口分配器(Window Assigner)API

分配器方法语法用途代码示例注意事项
TumblingEventTimeWindowsTumblingEventTimeWindows.of(Time.seconds(5))创建基于事件时间的滚动窗口stream.keyBy(key)
.window(TumblingEventTimeWindows.of(Time.seconds(5)))
推荐用于事件时间处理
TumblingProcessingTimeWindowsTumblingProcessingTimeWindows.of(Time.minutes(1))创建基于处理时间的滚动窗口.window(TumblingProcessingTimeWindows.of(Time.minutes(1)))不依赖数据时间戳
SlidingEventTimeWindowsSlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5))滑动窗口:长度10秒,每5秒滑动一次.window(SlidingEventTimeWindows.of(
Time.seconds(10), Time.seconds(5)))
第一个参数为窗口大小,第二个为滑动步长
SlidingProcessingTimeWindowsSlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(5))处理时间滑动窗口同上,替换为 ProcessingTime 版本同上
EventTimeSessionWindowsEventTimeSessionWindows.withGap(Time.minutes(10))事件时间会话窗口,固定 gap.window(EventTimeSessionWindows.withGap(
Time.minutes(10)))
gap 表示用户不活跃的最大间隔
DynamicEventTimeSessionWindowsDynamicEventTimeSessionWindows.withDynamicGap(...)动态 gap 会话窗口.window(DynamicEventTimeSessionWindows.withDynamicGap(
(element) -> element.getTimeout()))
根据元素动态设置 gap
GlobalWindowsGlobalWindows.create()将所有元素分配到单个全局窗口.window(GlobalWindows.create())必须自定义 Trigger,否则不输出
自定义 Window Assigner实现 WindowAssigner 接口自定义窗口逻辑class MyWindowAssigner extends WindowAssigner {...}高级用法,需理解 Flink 运行机制

5.3 窗口函数:ReduceFunction、AggregateFunction、ProcessWindowFunction

函数类型语法用途代码示例注意事项
ReduceFunctionReduceFunction<T>增量聚合,高效节省状态空间windowedStream.reduce((a, b) -> a.add(b));输入输出类型相同,仅支持简单聚合
AggregateFunctionAggregateFunction<IN, ACC, OUT>支持中间状态的增量聚合windowedStream.aggregate(new AverageAggregate());
其中 AverageAggregate 实现 createAccumulator, add, getResult, merge
类型可变,适合复杂聚合如平均值
ProcessWindowFunctionProcessWindowFunction<IN, OUT, KEY, W>全窗口函数,可访问上下文和元数据windowedStream.process(new ProcessWindowFunction<String, String, String, TimeWindow>() {
@Override
public void process(String key, Context ctx, Iterable<String> input, Collector<String> out) {
long windowStart = ctx.window().getStart();
out.collect("Window: " + windowStart + " Count: " + Iterators.size(input.iterator()));
}
})
可获取窗口元信息,但需缓存所有元素,资源开销大
增量 + 全窗口组合reduce/aggregate + ProcessWindowFunction增量预聚合 + 最终处理windowedStream.reduce(
(a, b) -> a + b,
new ProcessWindowFunction<Integer, String, String, TimeWindow>() {
@Override
public void process(...) { ... }
}
)
既高效又可访问上下文,推荐生产使用

5.4 窗口触发器(Trigger)与驱逐器(Evictor)

组件方法/类用途代码示例注意事项
Triggertrigger(Trigger)定义窗口何时触发计算.window(...).trigger(EventTimeTrigger.create())可自定义触发逻辑
EventTimeTriggerEventTimeTrigger.create()当 Watermark ≥ 窗口结束时间时触发默认用于事件时间窗口最常用触发器
ProcessingTimeTriggerProcessingTimeTrigger.create()基于处理时间触发用于处理时间窗口不处理乱序
CountTriggerCountTrigger.of(100)每积累100个元素触发一次.trigger(CountTrigger.of(100))可能提前触发,不保证完整性
PurgingTriggerPurgingTrigger.of(trigger)包装触发器,在触发后清空窗口内容.trigger(PurgingTrigger.of(EventTimeTrigger.create()))避免重复计算
ContinuousEventTimeTriggerContinuousEventTimeTrigger.of(Time.seconds(5))每隔5秒检查是否应触发.trigger(ContinuousEventTimeTrigger.of(Time.seconds(5)))支持周期性输出
Evictorevictor(Evictor)在触发前后移除部分元素.evictor(CountEvictor.of(10))影响性能,慎用
CountEvictorCountEvictor.of(10)保留最近10个元素,其余移除同上常用于滑动窗口优化
DeltaEvictorDeltaEvictor.of(threshold, comparator)移除与最近元素差异大的数据.evictor(DeltaEvictor.of(0.1, new MyComparator()))需实现比较逻辑

5.5 延迟数据处理:allowedLateness 与侧输出流(Side Output)

方法/组件语法用途代码示例注意事项
allowedLateness.allowedLateness(Time.minutes(1))允许迟到1分钟内的数据参与窗口计算windowedStream
.allowedLateness(Time.minutes(1))
.sum(1);
窗口保持打开,直到 Watermark > 窗口结束 + 延迟时间
再次触发-迟到数据到达时重新触发窗口计算结合 ProcessWindowFunction 可感知多次触发输出可能多次,需下游处理幂等
Side OutputOutputTag<T> + getSideOutput(tag)将迟到数据发送到侧输出流OutputTag<String> lateTag = new OutputTag<String>("late-data"){};
SingleOutputStreamOperator<...> result = stream
.window(...)
.allowedLateness(...)
.sideOutputLateData(lateTag);
DataStream<String> lateStream = result.getSideOutput(lateTag);
主流继续处理,迟到数据单独处理
sideOutputLateData.sideOutputLateData(OutputTag)指定迟到数据的输出标签见上例必须先调用 allowedLateness
处理侧输出result.getSideOutput(tag)获取侧输出流进行后续处理lateStream.addSink(new KafkaSink(...));可写入日志、降级存储或重试队列

第六章:状态管理与容错机制

6.1 状态类型:Keyed State 与 Operator State

状态类型说明适用场景注意事项
Keyed State与 key 相关的状态,每个 key 拥有独立状态keyBy 后的算子中维护计数、累计值等由 Flink 自动分区,支持大规模并行
Operator State与整个算子实例相关,不绑定 keySource 中记录偏移量、算子本地缓存需手动管理分片与恢复逻辑
原生状态 vs 托管状态托管状态由 Flink 管理(推荐),原生状态自行管理托管状态支持自动 checkpoint生产环境应使用托管状态

6.2 Keyed State API:ValueState、ListState、MapState 等

State 类型语法用途代码示例注意事项
ValueStateValueState<T>存储单个值,如最新温度private transient ValueState<Long> countState;
ValueStateDescriptor<Long> descriptor = new ValueStateDescriptor<>("count", TypeInformation.of(Long.class), 0L);
countState = getRuntimeContext().getState(descriptor);
countState.value(); // 读取
countState.update(10); // 写入
支持默认值,update 覆盖旧值
ListStateListState<T>存储列表,如用户点击序列ListState<String> clicks = context.getListState(
new ListStateDescriptor<>("clicks", Types.STRING)
);
for (String c : clicks.get()) { ... }
clicks.add("home");
可迭代,支持添加多个元素
ReducingStateReducingState<T>增量聚合状态,类似 reduceReducingState<Integer> sum = context.getReducingState(
new ReducingStateDescriptor<>("sum", (a, b) -> a + b, Integer.class)
);
高效,避免存储全部元素
AggregatingStateAggregatingState<IN, OUT>支持中间状态的聚合AggregatingState<Integer, Double> avg = context.getAggregatingState(
new AggregatingStateDescriptor<>("avg", new AvgFunction(), Types.DOUBLE)
);
类似 AggregateFunction
MapStateMapState<K, V>存储键值对映射MapState<String, Integer> wordCount = context.getMapState(
new MapStateDescriptor<>("wc", Types.STRING, Types.INT)
);
wordCount.put("hello", 1);
Integer c = wordCount.get("hello");
支持 put/get/contains/remove
状态清除state.clear()清除当前 key 的状态countState.clear();避免状态无限增长,建议在窗口结束或超时时清理

6.3 状态后端(State Backend):Memory、Fs、RocksDB

后端类型语法用途注意事项
MemoryStateBackendnew MemoryStateBackend()状态存储在 JVM 堆内存适用于小状态、本地测试
FsStateBackendnew FsStateBackend("file:///path", true)状态快照写入文件系统(HDFS/S3),堆内或堆外中等状态,生产常用
RocksDBStateBackendnew RocksDBStateBackend("file:///path")状态存储在本地 RocksDB,支持超大状态大状态、高可用场景推荐
配置方式env.setStateBackend(new RocksDBStateBackend(...))设置状态后端必须在创建状态前配置

6.4 Checkpointing 机制与配置

方法/配置语法用途代码示例注意事项
enableCheckpointingenv.enableCheckpointing(long interval)开启 Checkpoint,每 interval 毫秒执行一次env.enableCheckpointing(5000); // 5秒一次是容错基础
setCheckpointingMode.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE)设置一致性语义默认 EXACTLY_ONCE,可选 AT_LEAST_ONCE生产环境用 exactly-once
setCheckpointTimeout.setCheckpointTimeout(long timeout)Checkpoint 必须在超时内完成.setCheckpointTimeout(10000)超时则失败,不影响作业
setMinPauseBetweenCheckpoints.setMinPauseBetweenCheckpoints(long pause)两次 Checkpoint 最小间隔.setMinPauseBetweenCheckpoints(2000)避免频繁影响性能
setMaxConcurrentCheckpoints.setMaxConcurrentCheckpoints(1)最大并发 Checkpoint 数通常设为1,避免资源竞争可设为2以提高吞吐
enableExternalizedCheckpoints.enableExternalizedCheckpoints(CleanupMode)外部化 Checkpoint,重启不删除.enableExternalizedCheckpoints(
ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION)
便于手动恢复和 Savepoint 替代

6.5 Savepoint 的使用与应用升级

操作说明命令/代码示例注意事项
触发 Savepoint通过 CLI 手动生成快照bin/flink savepoint <jobId> hdfs:///flink/savepointsjobId 可从 Web UI 获取
从 Savepoint 启动使用 Savepoint 恢复作业bin/flink run -s hdfs:///flink/savepoints/savepoint-abc123 ...支持修改并行度、修复逻辑后重启
应用升级修改代码后从 Savepoint 重启先 cancel with savepoint,再 run -s算子 UID 不变才能正确恢复状态
设置算子 UID保证状态映射一致mapFunc.uid("mapper-1").setParallelism(4)推荐显式设置 UID,避免重构导致不匹配
兼容性要求状态结构变更限制字段增删需谨慎,建议向后兼容使用 Avro/Protobuf 更安全
清理策略管理过期 Savepoint手动删除 HDFS 文件或脚本定期清理Savepoint 不自动过期,需运维介入

7.1 Table API 与 SQL 简介

概念说明用途注意事项
Table API嵌入式 DSL,支持链式调用,类型安全适用于 Java/Scala 程序员,便于集成语法接近 SQL,但为编程式 API
Flink SQL支持标准 ANSI SQL 的查询语言快速开发、易读易维护,支持流批统一生产环境推荐用于 ETL、聚合场景
统一 APITable API 和 SQL 底层共用同一优化器(Calcite)两种方式可混合使用,互相转换通过 tableEnv.sqlQuery()table.toDataStream() 转换
流式 SQL支持在无限流上执行 SQL 查询实现 CEP、实时聚合、关联分析需定义时间属性和 Watermark
Blink PlannerFlink 默认的查询优化器提供更好的性能和功能支持无需显式设置,自动启用

7.2 创建 TableEnvironment

方法语法用途代码示例注意事项
StreamTableEnvironment.create()StreamTableEnvironment.create(env)创建流式表环境,关联 DataStreamStreamExecutionEnvironment dsEnv = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(dsEnv);
推荐用于流处理应用
BatchTableEnvironment.create()BatchTableEnvironment.create(env)创建批处理表环境(已弃用)新版本统一使用 StreamTableEnvironment流批统一后不再区分
EnvironmentSettingsEnvironmentSettings.newInstance().build()配置表环境设置EnvironmentSettings settings = EnvironmentSettings.newInstance()
.inStreamingMode()
.build();
TableEnvironment tableEnv = TableEnvironment.create(settings);
支持流/批模式选择,更灵活
设置并行度tableEnv.getConfig().set...配置表执行参数tableEnv.getConfig().getConfiguration()
.set(ExecutionCheckpointingOptions.CHECKPOINTING_INTERVAL, Duration.ofSeconds(10));
可设置 Checkpoint、状态后端等

7.3 注册表源(Kafka、文件等)与表 Sink

操作语法用途代码示例注意事项
createTemporaryTabletableEnv.createTemporaryTable("name", TableDescriptor)通过 Table API 注册临时表tableEnv.createTemporaryTable("kafka_source",
TableDescriptor.forConnector("kafka")
.option("topic", "input")
.option("properties.bootstrap.servers", "localhost:9092")
.format("json")
.option("scan.startup.mode", "earliest-offset")
.build());
推荐方式,声明式配置
executeSqltableEnv.executeSql("CREATE TABLE ...")使用 SQL DDL 创建表tableEnv.executeSql(
"CREATE TABLE file_sink (" +
" name STRING," +
" age INT" +
") WITH (" +
" 'connector' = 'filesystem'," +
" 'path' = '/tmp/output'," +
" 'format' = 'json'" +
")"
);
与 Hive 兼容,适合脚本化部署
Kafka 连接器.forConnector("kafka")读写 Kafka 主题见上例需引入 flink-sql-connector-kafka
File System 连接器.forConnector("filesystem")读写本地/HDFS 文件.option("path", "/data/input")支持有界和无界流
JDBC 连接器.forConnector("jdbc")连接关系型数据库.option("url", "jdbc:mysql://localhost:3306/test")
.option("table-name", "users")
仅支持作为 Sink,Source 需自定义
Upsert Kafka.forConnector("upsert-kafka")支持更新/删除的 Kafka 写入用于维表关联后写回需定义 PRIMARY KEY

7.4 执行 SQL 查询:SELECT、GROUP BY、JOIN(流式 Join)

查询类型语法用途代码示例注意事项
SELECT 投影SELECT col1, col2 FROM table字段选择与计算Table result = tableEnv.sqlQuery("SELECT user, price * 2 AS double_price FROM orders");支持表达式、函数
WHERE 过滤WHERE condition条件过滤tableEnv.sqlQuery("SELECT * FROM clicks WHERE page = 'home'");支持复杂布尔逻辑
GROUP BY 聚合GROUP BY key分组聚合tableEnv.sqlQuery("SELECT user, COUNT(*) FROM clicks GROUP BY user");流式聚合需定义时间窗口
窗口聚合结合 TUMBLE/HOP/SESSION时间窗口统计见 7.5 节必须基于时间属性
INNER JOINJOIN ... ON ...内连接,仅输出匹配行SELECT a.user, b.addr FROM clicks a JOIN users b ON a.user = b.user;流流 Join 需定义时间窗口
LEFT JOIN / RIGHT JOINLEFT JOIN ... ON ...左/右外连接SELECT a.user, b.addr FROM clicks a LEFT JOIN users b ON a.user = b.user;流式 LEFT JOIN 支持延迟右表更新
Temporal JoinFOR SYSTEM_TIME AS OF时态 Join(维表关联)SELECT o.amount, r.rate FROM orders o
JOIN rates FOR SYSTEM_TIME AS OF o.proctime AS r
ON o.currency = r.currency;
关联变化的历史维度数据
Interval JoinJOIN ... ON key AND t1.time BETWEEN ...区间 Join,限制时间范围SELECT * FROM orders o, shipments s
WHERE o.id = s.order_id
AND o.etime BETWEEN s.etime - INTERVAL '5' MINUTE AND s.etime;
自动清理状态,避免内存泄漏

7.5 时间属性在 SQL 中的定义(TUMBLE、HOP、SESSION)

时间函数语法用途代码示例注意事项
PROCTIME()PROCTIME() AS proc_time定义处理时间属性CREATE TABLE clicks (
user STRING,
url STRING,
pt AS PROCTIME()
) WITH (...);
自动生成,无需数据字段
ROWTIMEevent_time AS ROWTIME()从字段提取事件时间CREATE TABLE events (
ts BIGINT,
et AS TO_TIMESTAMP(FROM_UNIXTIME(ts)),
WATERMARK FOR et AS et - INTERVAL '5' SECOND
) WITH (...);
必须配合 WATERMARK 定义
WATERMARK FORWATERMARK FOR rowtime_col AS ...定义 Watermark 生成策略WATERMARK FOR et AS et - INTERVAL '10' SECOND延迟容忍设置
TUMBLETUMBLE(rowtime, INTERVAL '1' MINUTE)滚动窗口SELECT TUMBLE_START(et, INTERVAL '1' MINUTE), COUNT(*)
FROM events
GROUP BY TUMBLE(et, INTERVAL '1' MINUTE);
窗口开始/结束时间可提取
HOPHOP(rowtime, INTERVAL '5' SECONDS, INTERVAL '1' MINUTE)滑动窗口:每5秒滑动,窗口长1分钟SELECT HOP_START(et, INTERVAL '5' SECOND, INTERVAL '1' MINUTE), SUM(price)
FROM orders
GROUP BY HOP(et, INTERVAL '5' SECOND, INTERVAL '1' MINUTE);
第一个为滑动步长,第二个为窗口大小
SESSIONSESSION(rowtime, INTERVAL '10' MINUTES)会话窗口SELECT SESSION_START(et, INTERVAL '10' MINUTE), COUNT(*)
FROM clicks
GROUP BY SESSION(et, INTERVAL '10' MINUTE);
自动合并邻近会话

第八章:连接外部系统(Connectors)

8.1 Kafka Connector:读写 Kafka 数据

配置项语法用途代码示例注意事项
connector'connector' = 'kafka'指定 Kafka 连接器.option("connector", "kafka")必须
topic / topic-pattern'topic' = 'orders''topic-pattern' = 'test.*'指定单个或多个主题.option("topic", "user-behavior")支持正则匹配
properties.bootstrap.servers'properties.bootstrap.servers' = 'localhost:9092'Kafka 集群地址.option("properties.bootstrap.servers", "kafka1:9092,kafka2:9092")必须
format'format' = 'json''csv'数据序列化格式.option("format", "json")需引入对应 format 依赖
scan.startup.mode'scan.startup.mode' = 'earliest-offset'消费起始位置可选:earliest-offset, latest-offset, group-offsets, timestamptimestamp 需配合 scan.startup.timestamp-millis
sink.partitioner'sink.partitioner' = 'round-robin'写入分区策略可选:fixed, round-robin, custom默认按 key 分区
事务写入支持 exactly-once通过两阶段提交保障一致性默认启用,sink.delivery-guarantee = 'exactly-once'需 Kafka 0.11+

8.2 File System Connector(支持滚动文件)

配置项语法用途代码示例注意事项
connector'connector' = 'filesystem'使用文件系统连接器.option("connector", "filesystem")必须
path'path' = '/tmp/output'输出目录路径.option("path", "hdfs://namenode:8020/flink/output")支持本地、HDFS、S3
format'format' = 'json'文件输出格式.option("format", "json")也支持 csv, avro, parquet
sink.rolling-policy配置滚动策略控制文件切分.option("sink.rolling-policy.rollover-interval", "60s")
.option("sink.rolling-policy.file-size", "1GB")
按时间和大小滚动
分区写入动态分区按字段分区存储CREATE TABLE logs (
level STRING,
msg STRING,
dt STRING
) PARTITIONED BY (dt) WITH (...);
类似 Hive 分区
有界流 vs 无界流writeAsText 仅支持有界无界流需使用 Streaming File Sink推荐使用 SQL Connector更灵活,支持 Exactly-Once

8.3 JDBC Connector:写入 MySQL、PostgreSQL 等数据库

配置项语法用途代码示例注意事项
connector'connector' = 'jdbc'使用 JDBC 连接器.option("connector", "jdbc")必须
url'url' = 'jdbc:mysql://localhost:3306/test'数据库连接 URL.option("url", "jdbc:postgresql://localhost:5432/mydb")必须
table-name'table-name' = 'users'目标表名.option("table-name", "orders")必须
username / password'username' = 'root', 'password' = '123'认证信息.option("username", "flink")
.option("password", "secret")
建议使用密钥管理
sink.buffer-flush.interval'sink.buffer-flush.interval' = '2s'缓冲刷新间隔控制写入频率避免频繁小批量写入
sink.buffer-flush.max-rows'sink.buffer-flush.max-rows' = '1000'最大缓存行数达到后触发 flush平衡延迟与吞吐
支持数据库MySQL、PostgreSQL、Oracle、SQL Server 等通用 JDBC 驱动需将驱动 JAR 放入 lib 目录不支持自动 DDL 创建表

8.4 Elasticsearch Connector

配置项语法用途代码示例注意事项
connector'connector' = 'elasticsearch-7'指定 ES 7+ 连接器.option("connector", "elasticsearch-7")6.x 使用 elasticsearch-6
hosts'hosts' = 'http://localhost:9200'ES 集群地址.option("hosts", "http://es1:9200,http://es2:9200")必须
index'index' = 'flink-data'写入的索引名.option("index", "logs-{yyyy-MM-dd}")支持动态索引名
document-type'document-type' = '_doc'文档类型(ES 7+ 可省略).option("document-type", "_doc")ES 7 后默认 _doc
sink.flush-on-checkpoint'sink.flush-on-checkpoint' = 'true'Checkpoint 时刷新缓冲保障 Exactly-Once推荐开启
key-delimiter'key-delimiter' = '$'复合 key 的分隔符用于 upsert 操作默认 _
支持操作upsert-mode支持 insert、upsert需定义主键PRIMARY KEY (id) NOT ENFORCED

8.5 自定义 Connector 开发

组件说明开发要点注意事项
SourceFunction / Source实现 SourceFunction(批)或 RichSourceFunction覆盖 run()cancel() 方法run() 中循环 collector.collect() 发送数据
InputFormat / InputSplitSource批处理源,支持并行分片实现 InputFormat 接口用于有界数据源如数据库全量读取
SinkFunction实现 SinkFunction覆盖 invoke() 方法写入外部系统简单同步写入
RichSinkFunction继承 RichSinkFunction支持 open() 初始化连接,close() 释放资源建议使用,便于资源管理
Checkpointing 支持实现 CheckpointedFunction管理算子状态,支持故障恢复snapshotState() 保存状态,initializeState() 恢复
异步写入使用 AsyncFunction提高吞吐,避免阻塞适用于数据库、HTTP 等慢速系统
Connector DDL 支持实现 DynamicTableSource / DynamicTableSink支持 SQL 方式使用自定义连接器高级用法,需注册到 Factory

第九章:性能调优与监控

9.1 并行度设置与任务链(Task Chaining)

概念/方法语法用途代码示例注意事项
setParallelismenv.setParallelism(4)设置全局并行度StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(8);
可被算子级设置覆盖
operator.setParallelismmapOp.setParallelism(2)设置单个算子并行度DataStream<String> stream = env.addSource(...)
.map(...).setParallelism(4);
优先级高于全局设置
算子链(Operator Chain)Flink 自动将可链接算子合并为 Task减少序列化与网络开销,提升性能默认开启仅发生在同一任务槽(Task Slot)内
disableChainingenv.disableOperatorChaining()关闭自动任务链env.disableOperatorChaining();用于调试或资源隔离
startNewChain.startNewChain()从当前算子开始新链stream.map(...).startNewChain().filter(...)强制断开前一个链
disableChaining().disableChaining()当前算子不参与链接mapFunc.disableChaining()该算子独立运行
算子链策略env.getConfig().set...配置链行为可通过配置调整链合并策略通常使用默认即可

9.2 背压(Backpressure)识别与处理

方法/工具说明用途注意事项
Web UI Backpressure Indicator任务背压状态显示(High/Medium/Low)快速识别瓶颈算子若某算子背压为 High,说明其处理速度慢于上游
Metrics: Input/Output Queuemetrics.inputQueueLength, outputQueueLength监控任务输入输出缓冲区持续高队列长度表明处理延迟
Flink Web UI 吞吐图显示每秒输入/输出记录数分析数据流动瓶颈对比上下游吞吐,定位慢算子
降低并行度差异避免 map(16) -> keyBy -> reduce(2)防止下游成为瓶颈调整 reduce 并行度以匹配上游
增加资源增加 TaskManager 槽位或并行度缓解计算压力需评估数据倾斜问题
异步 I/O使用 AsyncFunction避免阻塞式外部调用导致背压见 10.2 节
检查状态访问避免在 ProcessFunction 中同步访问大状态减少处理延迟考虑使用 RocksDB 或异步查询

9.3 内存管理与配置调优

配置项语法/参数用途注意事项
TaskManager 内存模型taskmanager.memory.process.size总内存大小(如 4g)包含 JVM 堆、堆外、Metaspace、JVM 开销
Managed Memorytaskmanager.memory.managed.fraction.managed.size用于排序、JOIN、Caching 的托管内存批处理中尤为重要,流式可调小
Network Memorytaskmanager.memory.network.fraction网络缓冲区内存默认合理,高并发可适当调大
JVM 堆内存taskmanager.memory.task.heap.sizeTask 堆内存通常由总内存自动分配
Off-heap Memorytaskmanager.memory.task.off-heap: true使用堆外内存减少 GC 压力,适用于大状态
RocksDB 内存state.backend.rocksdb.memory.managed: true启用 RocksDB 托管内存自动管理内存,避免 OOM
配置方式flink-conf.yaml 或命令行 -D生产环境推荐使用配置文件-Dtaskmanager.memory.process.size=4g

9.4 指标(Metrics)系统:监控吞吐、延迟等

指标类型指标名称用途获取方式注意事项
CounternumRecordsIn, numRecordsOut记录输入/输出条数getRuntimeContext().getMetricGroup().counter("myCounter")可自定义计数器
GaugecurrentInputWatermark当前水位线值实时反映延迟通过 Web UI 查看
HistograminputDataSize数据大小分布分析数据量波动需启用 metrics.reporter
MeteroutPerSec每秒输出速率(吞吐)getRuntimeContext().getMetricGroup().meter("throughput", new MeterView(60))评估系统性能
注册自定义指标Counter counter = getRuntimeContext().getMetricGroup().counter("elements")监控业务指标counter.inc();open() 中初始化
报告器配置metrics.reporter.jmx.class: org.apache.flink.metrics.jmx.JMXReporter将指标暴露给 JMX也可配置 Prometheus、Graphite便于集成监控系统

9.5 Web UI 使用与日志分析

功能说明用途注意事项
Job Overview显示所有作业状态(Running/Finished/Failed)监控作业生命周期点击进入详情页
Task Details查看算子并行度、吞吐、延迟、反压定位性能瓶颈关注 Input/Output Records, Latency
Checkpointing显示 Checkpoint 成功率、耗时、大小评估容错性能失败或超时需排查网络或状态过大
Logs提供 TaskManager 和 JobManager 日志链接排查异常与错误日志路径通常为 log/flink-*.log
查看异常堆栈Failed 作业的异常信息定位代码错误如 NullPointerException, SerializationException
Metrics 图表实时展示自定义和系统指标可视化监控支持导出为 PNG
Timeline作业执行时间线分析调度延迟从提交到运行的时间

10.1 容错与 Exactly-Once 语义实现

组件说明实现方式注意事项
Checkpointing分布式快照机制周期性保存算子状态需启用并配置间隔、超时
TwoPhaseCommitSinkFunction两阶段提交 Sink实现端到端 Exactly-OnceKafka、JDBC、Elasticsearch 等均基于此
Barrier 对齐In-flight 数据与 Checkpoint 对齐保证状态一致性仅在 Exactly-Once 模式下需要
Source 支持可重放的 Source(如 Kafka)故障后从 Checkpoint 位点恢复文件源需保证不可变
Sink 支持支持事务或幂等写入Kafka 使用事务,数据库使用 XA避免重复写入
Watermark 一致性Watermark 随状态保存恢复后继续推进时间保证事件时间语义
配置建议env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
启用 Exactly-Once网络稳定、状态后端可靠

10.2 异步 I/O(Async I/O)访问外部系统

方法语法用途代码示例注意事项
AsyncFunctionAsyncFunction<IN, OUT>异步访问数据库、HTTP、Redispublic class AsyncDatabaseRequest extends AsyncFunction<String, String> {
@Override
public void asyncInvoke(String input, ResultFuture<String> resultFuture) {
httpClient.get(input, response -> resultFuture.complete(Collections.singletonList(response)));
}
}.returns(Types.STRING);
必须调用 resultFuture.complete()
AsyncDataStream.orderedWaitAsyncDataStream.orderedWait(stream, ...)保持输入顺序AsyncDataStream.orderedWait(stream, func, 1000, TimeUnit.MILLISECONDS, 100);延迟由最慢请求决定
AsyncDataStream.unorderedWaitAsyncDataStream.unorderedWait(...)不保证顺序,低延迟AsyncDataStream.unorderedWait(stream, func, 1000, 100);输出可能乱序,适合日志处理
并发限制参数 capacity控制最大并发请求数避免压垮外部系统建议根据目标系统 QPS 设置
超时处理设置超时时间防止异步请求挂起resultFuture.completeExceptionally(new TimeoutException());需处理失败情况

10.3 广播状态(Broadcast State)模式

概念语法用途代码示例注意事项
BroadcastStreambroadcastStream.broadcast()创建广播流BroadcastStream<Config> broadcast = configStream.broadcast(configStateDescriptor);通常用于配置、规则
KeyedStream.connect(BroadcastStream)stream.connect(broadcast)连接主数据流与广播流stream.connect(broadcast).process(new BroadcastProcessFunction<...>() {...})主流可 Keyed,广播流无 Key
MapStateDescriptornew MapStateDescriptor<>("config", String.class, Integer.class)定义广播状态结构传递给 broadcast() 方法状态对所有并行实例共享
BroadcastProcessFunctionBroadcastProcessFunction<IN, BC, OUT>处理广播状态public void processBroadcastElement(Config value, Context ctx, Collector<OUT> out) {
ctx.getBroadcastState(configDesc).put("threshold", value.getThreshold());
}
可更新状态
ProcessElementprocessElement()处理主流数据Integer threshold = ctx.getBroadcastState(configDesc).get("threshold");可读取最新广播状态
状态一致性所有并行实例状态一致用于规则引擎、动态配置不支持 Keyed Broadcast State-

10.4 用户自定义函数(UDF)与函数类

类型语法用途代码示例注意事项
MapFunctionMapFunction<IN, OUT>一对一转换stream.map(new MapFunction<String, Integer>() {
public Integer map(String s) { return s.length(); }
}).returns(Types.INT);
简单转换
FilterFunctionFilterFunction<T>过滤数据stream.filter(new FilterFunction<String>() {
public boolean filter(String s) { return s.startsWith("error"); }
});
返回 true 保留
RichFunctionRichMapFunction, RichFlatMapFunction支持 open(), close(), getRuntimeContext()public class RichCounter extends RichMapFunction<String, Tuple2<String, Long>> {
private transient Counter counter;
public void open(Configuration c) {
counter = getRuntimeContext().getMetricGroup().counter("lines");
}
}
推荐使用
UDF 序列化实现 java.io.Serializable保证函数可序列化Lambda 表达式需注意捕获变量避免引用不可序列化对象
函数类打包flink run -c MyJob jarfile确保 UDF 在 JAR 中使用 maven-shade-plugin 打包依赖避免 ClassNotFoundException

10.5 生产环境部署与升级策略

策略说明操作方式注意事项
高可用部署防止单点故障配置 ZooKeeper 或 Kubernetes HAJobManager 可故障转移
Savepoint 升级修改逻辑后从 Savepoint 恢复bin/flink cancel -s hdfs://savepoint jobId
bin/flink run -s hdfs://savepoint newJobJar
保证算子 UID 不变
并行度调整扩容/缩容从 Savepoint 启动时指定新并行度状态可重新分配
版本升级Flink 大版本升级(如 1.15 -> 1.18)先测试兼容性,逐步灰度检查 API 变更与废弃项
灰度发布先上线部分任务使用不同 Job 名称或集群验证稳定后全量切换
监控告警集成 Prometheus + Grafana + Alertmanager监控 Checkpoint、背压、延迟设置阈值告警
资源隔离多作业分集群或命名空间Kubernetes 命名空间或 YARN 队列避免相互影响