Article
第一章:Flink 入门与核心概念
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(已弃用)。 |
1.2 Flink 架构概览:JobManager、TaskManager、Client
| 组件名称 | 说明 | 注意事项 |
|---|---|---|
| JobManager | 负责调度任务、协调检查点(Checkpoint)、管理元数据。集群中可有多个,但只有一个为 Leader。 | 高可用模式下需配合 ZooKeeper 或 Kubernetes 实现故障转移。 |
| TaskManager | 实际执行任务的工作节点,管理内存、网络、任务线程。多个 TaskManager 形成资源池。 | 每个 Task Slot 执行一个或多个算子子任务,Slot 数量影响并行度上限。 |
| Client | 提交作业的客户端,负责解析程序、生成 JobGraph 并提交给 JobManager。 | Client 可运行在本地或远程,提交后可断开(detached 模式)。 |
| Task Slot | TaskManager 中的资源单位,每个 Slot 独立执行一个任务链(task chain)。 | Slot 数量由内存决定,不隔离 CPU,建议根据资源需求合理配置。 |
| JobGraph / ExecutionGraph | Client 将程序转换为 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 Unbounded | Flink 将有界流视为无界流的特例,统一调度和执行。 | 有界流作业执行完毕后自动终止,无界流持续运行。 |
| Transformation | 数据转换操作,如 map、keyBy、window 等,构建 DAG 图。 | 转换是懒加载的,只有调用 execute() 才触发执行。 |
1.4 Flink 应用程序的执行流程
| 阶段 | 说明 | 注意事项 |
|---|---|---|
| 1. 编写程序 | 用户使用 DataStream API 或 Table API 编写逻辑。 | 程序需定义执行环境(StreamExecutionEnvironment)。 |
| 2. 创建 StreamGraph | 客户端根据 API 调用生成逻辑图(StreamGraph)。 | 图中节点为算子,边为数据流。 |
| 3. 生成 JobGraph | 客户端优化 StreamGraph(如任务链合并),生成 JobGraph 提交至 JobManager。 | JobGraph 是调度的基本单位。 |
| 4. 生成 ExecutionGraph | JobManager 将 JobGraph 转为 ExecutionGraph,分配任务到 TaskManager。 | ExecutionGraph 包含任务的并行实例(subtask)。 |
| 5. 调度与执行 | TaskManager 启动任务,开始消费数据并执行计算。 | 支持 Checkpoint 机制保障容错。 |
| 6. 结果输出 | 数据通过 Sink 输出到外部系统(如 Kafka、文件、数据库)。 | Sink 可配置异步或同步写入。 |
1.5 Flink 集群部署模式:Local、Standalone、YARN、Kubernetes
| 部署模式 | 说明 | 注意事项 |
|---|---|---|
| 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 或 11 | Flink 基于 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 用户注意路径分隔符。 |
2.2 使用 Maven 构建 Flink 项目
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 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 |
2.3 编写第一个 Flink 程序:WordCount(流式与批式)
| 程序类型 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| 批式 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(); | 自动识别运行环境,开发调试最常用 |
| createLocalEnvironment | StreamExecutionEnvironment.createLocalEnvironment() | 创建本地执行环境,指定并行度 | final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(2); | 用于本地测试,限制并行度为 2 |
| createRemoteEnvironment | StreamExecutionEnvironment.createRemoteEnvironment(String host, int port, String... jarFiles) | 连接远程 Flink 集群提交作业 | final StreamExecutionEnvironment env = StreamExecutionEnvironment.createRemoteEnvironment( "jobmanager", 8081, "path/to/job.jar"); | 需指定 JobManager 地址和 JAR 路径 |
| setParallelism | env.setParallelism(int parallelism) | 设置整个程序的默认并行度 | env.setParallelism(4); | 可被算子级并行度覆盖 |
| disableOperatorChaining | env.disableOperatorChaining() | 禁用任务链(chaining),每个算子独立调度 | env.disableOperatorChaining(); | 调试时便于观察子任务 |
| configure | env.getConfig().set... | 配置运行时参数,如类型序列化、全局作业参数 | env.getConfig().setAutoWatermarkInterval(200); | 常用于设置水印生成间隔 |
3.2 数据源(Source):从集合、文件、Socket、Kafka 读取数据
| Source 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| fromElements | env.fromElements(T... elements) | 从 Java 元素创建数据流,用于测试 | DataStream<String> stream = env.fromElements("a", "b", "c"); | 数据量小,适合本地调试 |
| fromCollection | env.fromCollection(Collection<T> collection) | 从集合创建流 | List<String> data = Arrays.asList("x", "y");DataStream<String> stream = env.fromCollection(data); | 支持任意 Collection 类型 |
| readTextFile | env.readTextFile(String path) | 读取文本文件,每行作为一个元素 | DataStream<String> lines = env.readTextFile("/tmp/input.txt"); | 路径支持本地或 HDFS,文件必须有界 |
| socketTextStream | env.socketTextStream(String hostname, int port) | 从 Socket 读取文本流,按行分割 | DataStream<String> stream = env.socketTextStream("localhost", 9999); | 仅用于测试,生产环境不推荐 |
| addSource | env.addSource(SourceFunction<T> function) | 添加自定义数据源 | env.addSource(new FlinkKafkaConsumer<>( "topic", new SimpleStringSchema(), props)); | 通用接口,支持 Kafka、MySQL Binlog 等 |
| FlinkKafkaConsumer | new 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
| 转换操作 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| map | stream.map(MapFunction<T, R>) | 将每个元素转换为另一个元素 | DataStream<Integer> doubled = stream.map(x -> x * 2); | 返回一对一映射结果 |
| flatMap | stream.flatMap(FlatMapFunction<T, R>) | 将每个元素转换为零或多个元素 | stream.flatMap((String line, Collector<String> out) -> { for (String word : line.split(" ")) { out.collect(word); }}).returns(Types.STRING); | 常用于分词、展开嵌套结构 |
| filter | stream.filter(FilterFunction<T>) | 过滤不符合条件的元素 | DataStream<String> filtered = stream.filter(s -> !s.isEmpty()); | 返回布尔值,true 保留,false 丢弃 |
| keyBy | stream.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 |
| returns | flatMap(...).returns(TypeInformation) | 显式声明泛型类型信息 | .returns(Types.STRING) | 使用 Lambda 时防止类型擦除,必须调用 |
3.4 输出操作(Sink):打印、写入文件、Kafka、自定义 Sink
| Sink 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| print / printToErr | stream.print() / stream.printToErr() | 将数据打印到标准输出或错误流 | stream.print();stream.print("prefix: "); | 输出带并行子任务 ID,如 1> hello |
| writeAsText | stream.writeAsText(String path) | 将数据写入文本文件 | stream.writeAsText("/tmp/output"); | 仅用于有界流,生产环境不推荐 |
| addSink | stream.addSink(SinkFunction<T> function) | 添加自定义 Sink | stream.addSink(new FlinkKafkaProducer<>( "topic", new SimpleStringSchema(), props)); | 通用接口,支持多种外部系统 |
| FlinkKafkaProducer | new 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 依赖 |
| 自定义 Sink | class 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 执行程序与设置并行度
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| execute | env.execute(String jobName) | 触发程序执行,启动作业 | env.execute("My Flink Job"); | 必须调用,否则程序不运行 |
| setParallelism | env.setParallelism(int p) | 设置全局并行度 | env.setParallelism(4); | 可被算子级并行度覆盖 |
| setParallelism(算子级) | stream.map(...).setParallelism(2) | 为特定算子设置并行度 | stream.map(x -> x + 1).setParallelism(2); | 优先级高于全局设置 |
| getMaxParallelism | env.setMaxParallelism(int max) | 设置最大并行度,影响状态后端分片 | env.setMaxParallelism(128); | 影响 Checkpoint 和状态重平衡 |
| disableOperatorChaining | env.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 使用
| 方法/类 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| WatermarkStrategy | WatermarkStrategy.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 类型时间戳(毫秒) |
| assignTimestampsAndWatermarks | stream.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
| 分配器方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| TumblingEventTimeWindows | TumblingEventTimeWindows.of(Time.seconds(5)) | 创建基于事件时间的滚动窗口 | stream.keyBy(key).window(TumblingEventTimeWindows.of(Time.seconds(5))) | 推荐用于事件时间处理 |
| TumblingProcessingTimeWindows | TumblingProcessingTimeWindows.of(Time.minutes(1)) | 创建基于处理时间的滚动窗口 | .window(TumblingProcessingTimeWindows.of(Time.minutes(1))) | 不依赖数据时间戳 |
| SlidingEventTimeWindows | SlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5)) | 滑动窗口:长度10秒,每5秒滑动一次 | .window(SlidingEventTimeWindows.of( Time.seconds(10), Time.seconds(5))) | 第一个参数为窗口大小,第二个为滑动步长 |
| SlidingProcessingTimeWindows | SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(5)) | 处理时间滑动窗口 | 同上,替换为 ProcessingTime 版本 | 同上 |
| EventTimeSessionWindows | EventTimeSessionWindows.withGap(Time.minutes(10)) | 事件时间会话窗口,固定 gap | .window(EventTimeSessionWindows.withGap( Time.minutes(10))) | gap 表示用户不活跃的最大间隔 |
| DynamicEventTimeSessionWindows | DynamicEventTimeSessionWindows.withDynamicGap(...) | 动态 gap 会话窗口 | .window(DynamicEventTimeSessionWindows.withDynamicGap( (element) -> element.getTimeout())) | 根据元素动态设置 gap |
| GlobalWindows | GlobalWindows.create() | 将所有元素分配到单个全局窗口 | .window(GlobalWindows.create()) | 必须自定义 Trigger,否则不输出 |
| 自定义 Window Assigner | 实现 WindowAssigner 接口 | 自定义窗口逻辑 | class MyWindowAssigner extends WindowAssigner {...} | 高级用法,需理解 Flink 运行机制 |
5.3 窗口函数:ReduceFunction、AggregateFunction、ProcessWindowFunction
| 函数类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| ReduceFunction | ReduceFunction<T> | 增量聚合,高效节省状态空间 | windowedStream.reduce((a, b) -> a.add(b)); | 输入输出类型相同,仅支持简单聚合 |
| AggregateFunction | AggregateFunction<IN, ACC, OUT> | 支持中间状态的增量聚合 | windowedStream.aggregate(new AverageAggregate());其中 AverageAggregate 实现 createAccumulator, add, getResult, merge | 类型可变,适合复杂聚合如平均值 |
| ProcessWindowFunction | ProcessWindowFunction<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)
| 组件 | 方法/类 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| Trigger | trigger(Trigger) | 定义窗口何时触发计算 | .window(...).trigger(EventTimeTrigger.create()) | 可自定义触发逻辑 |
| EventTimeTrigger | EventTimeTrigger.create() | 当 Watermark ≥ 窗口结束时间时触发 | 默认用于事件时间窗口 | 最常用触发器 |
| ProcessingTimeTrigger | ProcessingTimeTrigger.create() | 基于处理时间触发 | 用于处理时间窗口 | 不处理乱序 |
| CountTrigger | CountTrigger.of(100) | 每积累100个元素触发一次 | .trigger(CountTrigger.of(100)) | 可能提前触发,不保证完整性 |
| PurgingTrigger | PurgingTrigger.of(trigger) | 包装触发器,在触发后清空窗口内容 | .trigger(PurgingTrigger.of(EventTimeTrigger.create())) | 避免重复计算 |
| ContinuousEventTimeTrigger | ContinuousEventTimeTrigger.of(Time.seconds(5)) | 每隔5秒检查是否应触发 | .trigger(ContinuousEventTimeTrigger.of(Time.seconds(5))) | 支持周期性输出 |
| Evictor | evictor(Evictor) | 在触发前后移除部分元素 | .evictor(CountEvictor.of(10)) | 影响性能,慎用 |
| CountEvictor | CountEvictor.of(10) | 保留最近10个元素,其余移除 | 同上 | 常用于滑动窗口优化 |
| DeltaEvictor | DeltaEvictor.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 Output | OutputTag<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 | 与整个算子实例相关,不绑定 key | Source 中记录偏移量、算子本地缓存 | 需手动管理分片与恢复逻辑 |
| 原生状态 vs 托管状态 | 托管状态由 Flink 管理(推荐),原生状态自行管理 | 托管状态支持自动 checkpoint | 生产环境应使用托管状态 |
6.2 Keyed State API:ValueState、ListState、MapState 等
| State 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| ValueState | ValueState<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 覆盖旧值 |
| ListState | ListState<T> | 存储列表,如用户点击序列 | ListState<String> clicks = context.getListState( new ListStateDescriptor<>("clicks", Types.STRING));for (String c : clicks.get()) { ... }clicks.add("home"); | 可迭代,支持添加多个元素 |
| ReducingState | ReducingState<T> | 增量聚合状态,类似 reduce | ReducingState<Integer> sum = context.getReducingState( new ReducingStateDescriptor<>("sum", (a, b) -> a + b, Integer.class)); | 高效,避免存储全部元素 |
| AggregatingState | AggregatingState<IN, OUT> | 支持中间状态的聚合 | AggregatingState<Integer, Double> avg = context.getAggregatingState( new AggregatingStateDescriptor<>("avg", new AvgFunction(), Types.DOUBLE)); | 类似 AggregateFunction |
| MapState | MapState<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
| 后端类型 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| MemoryStateBackend | new MemoryStateBackend() | 状态存储在 JVM 堆内存 | 适用于小状态、本地测试 |
| FsStateBackend | new FsStateBackend("file:///path", true) | 状态快照写入文件系统(HDFS/S3),堆内或堆外 | 中等状态,生产常用 |
| RocksDBStateBackend | new RocksDBStateBackend("file:///path") | 状态存储在本地 RocksDB,支持超大状态 | 大状态、高可用场景推荐 |
| 配置方式 | env.setStateBackend(new RocksDBStateBackend(...)) | 设置状态后端 | 必须在创建状态前配置 |
6.4 Checkpointing 机制与配置
| 方法/配置 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| enableCheckpointing | env.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/savepoints | jobId 可从 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 不自动过期,需运维介入 |
第七章:Flink SQL 与 Table API
7.1 Table API 与 SQL 简介
| 概念 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
| Table API | 嵌入式 DSL,支持链式调用,类型安全 | 适用于 Java/Scala 程序员,便于集成 | 语法接近 SQL,但为编程式 API |
| Flink SQL | 支持标准 ANSI SQL 的查询语言 | 快速开发、易读易维护,支持流批统一 | 生产环境推荐用于 ETL、聚合场景 |
| 统一 API | Table API 和 SQL 底层共用同一优化器(Calcite) | 两种方式可混合使用,互相转换 | 通过 tableEnv.sqlQuery() 或 table.toDataStream() 转换 |
| 流式 SQL | 支持在无限流上执行 SQL 查询 | 实现 CEP、实时聚合、关联分析 | 需定义时间属性和 Watermark |
| Blink Planner | Flink 默认的查询优化器 | 提供更好的性能和功能支持 | 无需显式设置,自动启用 |
7.2 创建 TableEnvironment
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| StreamTableEnvironment.create() | StreamTableEnvironment.create(env) | 创建流式表环境,关联 DataStream | StreamExecutionEnvironment dsEnv = StreamExecutionEnvironment.getExecutionEnvironment();StreamTableEnvironment tableEnv = StreamTableEnvironment.create(dsEnv); | 推荐用于流处理应用 |
| BatchTableEnvironment.create() | BatchTableEnvironment.create(env) | 创建批处理表环境(已弃用) | 新版本统一使用 StreamTableEnvironment | 流批统一后不再区分 |
| EnvironmentSettings | EnvironmentSettings.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
| 操作 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| createTemporaryTable | tableEnv.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()); | 推荐方式,声明式配置 |
| executeSql | tableEnv.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 JOIN | JOIN ... ON ... | 内连接,仅输出匹配行 | SELECT a.user, b.addr FROM clicks a JOIN users b ON a.user = b.user; | 流流 Join 需定义时间窗口 |
| LEFT JOIN / RIGHT JOIN | LEFT JOIN ... ON ... | 左/右外连接 | SELECT a.user, b.addr FROM clicks a LEFT JOIN users b ON a.user = b.user; | 流式 LEFT JOIN 支持延迟右表更新 |
| Temporal Join | FOR SYSTEM_TIME AS OF | 时态 Join(维表关联) | SELECT o.amount, r.rate FROM orders oJOIN rates FOR SYSTEM_TIME AS OF o.proctime AS rON o.currency = r.currency; | 关联变化的历史维度数据 |
| Interval Join | JOIN ... ON key AND t1.time BETWEEN ... | 区间 Join,限制时间范围 | SELECT * FROM orders o, shipments sWHERE o.id = s.order_idAND 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 (...); | 自动生成,无需数据字段 |
| ROWTIME | event_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 FOR | WATERMARK FOR rowtime_col AS ... | 定义 Watermark 生成策略 | WATERMARK FOR et AS et - INTERVAL '10' SECOND | 延迟容忍设置 |
| TUMBLE | TUMBLE(rowtime, INTERVAL '1' MINUTE) | 滚动窗口 | SELECT TUMBLE_START(et, INTERVAL '1' MINUTE), COUNT(*)FROM eventsGROUP BY TUMBLE(et, INTERVAL '1' MINUTE); | 窗口开始/结束时间可提取 |
| HOP | HOP(rowtime, INTERVAL '5' SECONDS, INTERVAL '1' MINUTE) | 滑动窗口:每5秒滑动,窗口长1分钟 | SELECT HOP_START(et, INTERVAL '5' SECOND, INTERVAL '1' MINUTE), SUM(price)FROM ordersGROUP BY HOP(et, INTERVAL '5' SECOND, INTERVAL '1' MINUTE); | 第一个为滑动步长,第二个为窗口大小 |
| SESSION | SESSION(rowtime, INTERVAL '10' MINUTES) | 会话窗口 | SELECT SESSION_START(et, INTERVAL '10' MINUTE), COUNT(*)FROM clicksGROUP 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, timestamp | timestamp 需配合 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)
| 概念/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| setParallelism | env.setParallelism(4) | 设置全局并行度 | StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(8); | 可被算子级设置覆盖 |
| operator.setParallelism | mapOp.setParallelism(2) | 设置单个算子并行度 | DataStream<String> stream = env.addSource(...) .map(...).setParallelism(4); | 优先级高于全局设置 |
| 算子链(Operator Chain) | Flink 自动将可链接算子合并为 Task | 减少序列化与网络开销,提升性能 | 默认开启 | 仅发生在同一任务槽(Task Slot)内 |
| disableChaining | env.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 Queue | metrics.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 Memory | taskmanager.memory.managed.fraction 或 .managed.size | 用于排序、JOIN、Caching 的托管内存 | 批处理中尤为重要,流式可调小 |
| Network Memory | taskmanager.memory.network.fraction | 网络缓冲区内存 | 默认合理,高并发可适当调大 |
| JVM 堆内存 | taskmanager.memory.task.heap.size | Task 堆内存 | 通常由总内存自动分配 |
| Off-heap Memory | taskmanager.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)系统:监控吞吐、延迟等
| 指标类型 | 指标名称 | 用途 | 获取方式 | 注意事项 |
|---|---|---|---|---|
| Counter | numRecordsIn, numRecordsOut | 记录输入/输出条数 | getRuntimeContext().getMetricGroup().counter("myCounter") | 可自定义计数器 |
| Gauge | currentInputWatermark | 当前水位线值 | 实时反映延迟 | 通过 Web UI 查看 |
| Histogram | inputDataSize | 数据大小分布 | 分析数据量波动 | 需启用 metrics.reporter |
| Meter | outPerSec | 每秒输出速率(吞吐) | 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 | 作业执行时间线 | 分析调度延迟 | 从提交到运行的时间 |
第十章:Flink 高级特性与最佳实践
10.1 容错与 Exactly-Once 语义实现
| 组件 | 说明 | 实现方式 | 注意事项 |
|---|---|---|---|
| Checkpointing | 分布式快照机制 | 周期性保存算子状态 | 需启用并配置间隔、超时 |
| TwoPhaseCommitSinkFunction | 两阶段提交 Sink | 实现端到端 Exactly-Once | Kafka、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)访问外部系统
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| AsyncFunction | AsyncFunction<IN, OUT> | 异步访问数据库、HTTP、Redis | public 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.orderedWait | AsyncDataStream.orderedWait(stream, ...) | 保持输入顺序 | AsyncDataStream.orderedWait(stream, func, 1000, TimeUnit.MILLISECONDS, 100); | 延迟由最慢请求决定 |
| AsyncDataStream.unorderedWait | AsyncDataStream.unorderedWait(...) | 不保证顺序,低延迟 | AsyncDataStream.unorderedWait(stream, func, 1000, 100); | 输出可能乱序,适合日志处理 |
| 并发限制 | 参数 capacity | 控制最大并发请求数 | 避免压垮外部系统 | 建议根据目标系统 QPS 设置 |
| 超时处理 | 设置超时时间 | 防止异步请求挂起 | resultFuture.completeExceptionally(new TimeoutException()); | 需处理失败情况 |
10.3 广播状态(Broadcast State)模式
| 概念 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| BroadcastStream | broadcastStream.broadcast() | 创建广播流 | BroadcastStream<Config> broadcast = configStream.broadcast(configStateDescriptor); | 通常用于配置、规则 |
| KeyedStream.connect(BroadcastStream) | stream.connect(broadcast) | 连接主数据流与广播流 | stream.connect(broadcast).process(new BroadcastProcessFunction<...>() {...}) | 主流可 Keyed,广播流无 Key |
| MapStateDescriptor | new MapStateDescriptor<>("config", String.class, Integer.class) | 定义广播状态结构 | 传递给 broadcast() 方法 | 状态对所有并行实例共享 |
| BroadcastProcessFunction | BroadcastProcessFunction<IN, BC, OUT> | 处理广播状态 | public void processBroadcastElement(Config value, Context ctx, Collector<OUT> out) { ctx.getBroadcastState(configDesc).put("threshold", value.getThreshold());} | 可更新状态 |
| ProcessElement | processElement() | 处理主流数据 | Integer threshold = ctx.getBroadcastState(configDesc).get("threshold"); | 可读取最新广播状态 |
| 状态一致性 | 所有并行实例状态一致 | 用于规则引擎、动态配置 | 不支持 Keyed Broadcast State | - |
10.4 用户自定义函数(UDF)与函数类
| 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| MapFunction | MapFunction<IN, OUT> | 一对一转换 | stream.map(new MapFunction<String, Integer>() { public Integer map(String s) { return s.length(); }}).returns(Types.INT); | 简单转换 |
| FilterFunction | FilterFunction<T> | 过滤数据 | stream.filter(new FilterFunction<String>() { public boolean filter(String s) { return s.startsWith("error"); }}); | 返回 true 保留 |
| RichFunction | RichMapFunction, 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 HA | JobManager 可故障转移 |
| Savepoint 升级 | 修改逻辑后从 Savepoint 恢复 | bin/flink cancel -s hdfs://savepoint jobIdbin/flink run -s hdfs://savepoint newJobJar | 保证算子 UID 不变 |
| 并行度调整 | 扩容/缩容 | 从 Savepoint 启动时指定新并行度 | 状态可重新分配 |
| 版本升级 | Flink 大版本升级(如 1.15 -> 1.18) | 先测试兼容性,逐步灰度 | 检查 API 变更与废弃项 |
| 灰度发布 | 先上线部分任务 | 使用不同 Job 名称或集群 | 验证稳定后全量切换 |
| 监控告警 | 集成 Prometheus + Grafana + Alertmanager | 监控 Checkpoint、背压、延迟 | 设置阈值告警 |
| 资源隔离 | 多作业分集群或命名空间 | Kubernetes 命名空间或 YARN 队列 | 避免相互影响 |