☰
Flink 1.16 发布说明深度解读:重试 Lookup Join、异步输出模式与网络超支缓冲等关键变更
2026/9/25 5:26:28 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

Apache Flink 1.16 是一次以"稳定性加固 + 易用性改进"为主线的版本。本篇以官方发布说明(docs/content/release-notes/flink-1.16.md)为主体,完整梳理从 1.15 升级到 1.16 时需要关注的配置、行为与依赖变更,并逐条对照当前仓库中的源码定义,帮助你判断哪些变更会直接影响现有作业的升级路径。读完后,你将能够回答三个问题:哪些配置项的默认行为在 1.16 发生了变化、哪些 API 被移除需要迁移、以及新增的 Retryable Lookup Join / 异步输出模式 / 超支缓冲(Overdraft Buffer)分别如何在实际作业中配置。

升级总览:先记住这些"默认值变化"

Flink 1.16 发布说明覆盖六大类变更:Clusters & Deployment、Table API & SQL、Connectors、Runtime & Coordination、Checkpoints、Python,外加依赖升级。其中需要特别注意的默认行为变化(即"不改任何配置、升级后行为也会变")有以下几处,建议在升级演练中重点回归:

  1. Hive Sink 在批模式下默认向 Hive Metastore 上报统计信息(FLINK-28883);
  2. 非对齐 Checkpoint(Unaligned Checkpoint)在反压场景下,每个 gate默认允许申请 5 个超支网络缓冲区(FLINK-26762),可能略微增加作业内存占用;
  3. Application 模式 + HA 开启时,JobID 不再固定为0000000000...,而是基于 cluster ID 生成(FLINK-19358);
  4. REST API 在组件未就绪时返回503 Service Unavailable而非 500 Internal Server Error(FLINK-25269);
  5. Avro 生成代码的命名空间改为org.apache.flink.avro.generated(FLINK-25962)。

Clusters & Deployment:jobmanager.sh 的 host/web-ui-port 参数弃用

发布说明指出(FLINK-28735):jobmanager.sh脚本中的host与web-ui-port命令行参数已被弃用,应改用对应的动态属性(dynamic properties)来指定,例如通过-Djobmanager.bind-host=...、-Drest.address=...、-Drest.port=...等 option 的方式传入。

迁移建议:在启动脚本中把原先写死在jobmanager.sh命令行的 host/port,统一挪到config.yaml(或flink-conf.yaml)与-D参数中管理,这样同一份集群配置可以在 standalone、YARN 等部署形态间复用。升级时若仍使用旧参数,1.16 会给出弃用告警,属于"软迁移"窗口期,但后续版本大概率会彻底移除。

Table API & SQL

移除字符串表达式 DSL(String Expression DSL)

发布说明(FLINK-26704)确认:从 Java/Scala/Python 三种 Table API 中移除了此前已弃用的 String 表达式 DSL,即形如table.select("$id + 1", "lower($name)")的字符串写法。

迁移方式:改用类型安全的表达式 API。Java 侧使用Expression构建器:

Expression idPlusOne = $id().plus(1); Expression lowerName = Functions.lower($name()); table.select(idPlusOne, lowerName);

Python(PyFlink)侧使用表达式 API 的函数式写法:

from pyflink.table import expressions as expr table.select(expr.col("id") + 1, expr.expr("lower(%s)" % "name").if_supported()) # 或更推荐: table.select(expr.col("id") + 1, expr.lower(expr.col("name")))

由于这是直接移除而非弃用,所有仍在使用字符串 DSL 的作业会在 1.16 上直接编译失败,升级前必须先完成改写。

新增:Retryable Lookup Join 解决外部维表更新延迟

发布说明(FLINK-28779)介绍了 1.16 的一个亮点能力:为同步/异步 Lookup Join 增加可重试查询(retryable lookup join),用来解决外部维表(如缓存、数据库)更新延迟导致的"查不到维度"问题。

配置方式:通过 HINT 指定重试谓词与重试策略

在仓库源码中,重试相关 HINT 键定义在 LookupJoinHintOptions.java 中,计划侧解析逻辑位于 FlinkHintStrategies.java 与 LookupJoinUtil.java。从源码结构看,重试配置支持四个键:

  • retry-predicate:判定失败类型(何种查询失败才需要重试);
  • retry-strategy:重试策略,支持固定间隔(FIXED_DELAY)与指数退避(EXPONENTIAL_DELAY)两类;
  • fixed-delay:固定延迟时间(如1s);
  • max-attempts:最大重试次数。

SQL 中的典型用法(通过/*+ LOOKUP(...) */HINT 注入到具体的 Lookup Join 节点上):

SELECT /*+ LOOKUP( 'table'='dim', 'retry-predicate'='EXCEPTION_AS_FAILURE', 'retry-strategy'='FIXED_DELAY', 'fixed-delay'='1s', 'max-attempts'='3' ) */ d.name, o.* FROM orders AS o JOIN dim FOR SYSTEM_TIME AS OF o.proc_time AS d ON o.id = d.id;

这些重试参数最终会被序列化进逻辑计划:执行节点 CommonExecLookupJoin.java 中定义了retryOptions字段(FIELD_NAME_RETRY_OPTIONS),流式执行节点 StreamExecLookupJoin.java 在运行时接收该字段并驱动重试。可以推断:retryOptions随计划 JSON 一起序列化,因此重试行为属于"编译期确定"的配置,升级/恢复 Savepoint 后行为保持一致。

新增:table.exec.async-lookup.output-mode让异步查询输出模式可配置

发布说明(FLINK-27622)建议:当结果不需要严格保序时,把新选项table.exec.async-lookup.output-mode设置为ALLOW_UNORDERED,在 append-only 流上可显著降低延迟。

该选项在仓库中的定义位于 ExecutionConfigOptions.java:

public static final ConfigOption<AsyncOutputMode> TABLE_EXEC_ASYNC_LOOKUP_OUTPUT_MODE = key("table.exec.async-lookup.output-mode") .enumType(AsyncOutputMode.class) .defaultValue(AsyncOutputMode.ORDERED) .withDescription( "Output mode for asynchronous operations which will convert to " + "AsyncDataStream.OutputMode, ORDERED by default. If set to " + "ALLOW_UNORDERED, will attempt to use AsyncDataStream.OutputMode.UNORDERED " + "when it does not affect the correctness of the result, otherwise " + "ORDERED will be still used.");

三点关键信息:

  1. 默认值是ORDERED,即 1.16 默认行为不变,只有显式配置才可能走无序模式;
  2. 枚举取值包含ORDERED、ALLOW_UNORDERED(以及显式的UNORDERED),底层映射到 DataStream API 的AsyncDataStream.OutputMode;
  3. 描述中明确了安全边界:设为ALLOW_UNORDERED时,只有在不影响结果正确性的情况下才真正切换为 UNORDERED,否则仍回退到 ORDERED——这是该选项比直接用 DataStream 异步模式更"保守"的设计。

与异步 Lookup Join 相关的完整配置族(同文件中定义)还包括:

配置项默认值说明
table.exec.async-lookup.output-modeORDERED异步查询输出模式,1.16 新增
table.exec.async-lookup.buffer-capacity100异步 I/O 最大并发数
table.exec.async-lookup.timeout3 min异步操作完成超时

加固:非确定性更新在 changelog 链路中的正确性检测

发布说明(FLINK-27849)提到:对于复杂流作业,现在可以在运行前检测并提示changelog 管线中存在的非确定性更新(non-deterministic updates)所导致的潜在正确性问题。这是一个"编译期告警"性质的加固:当优化器发现算子链路上存在可能被重算(retraction/retry)的非确定性计算(如非确定性的 UDF、可能改变结果集顺序的逻辑)时,会在作业部署前给出警告,避免脏数据静默流入下游。升级 1.16 后建议关注作业提交时的告警日志,对命中告警的链路做拆分或物化中间结果处理。

Connectors

Hive Sink:批模式下默认向 Metastore 上报统计信息

发布说明(FLINK-28883):批模式下,Hive Sink 现在默认会为写出的表和分区向 Hive Metastore 上报统计信息;文件数很多时该过程可能耗时较长,可通过table.exec.hive.sink.statistic-auto-gather.enable=false关闭。

在仓库中该选项定义于 HiveOptions.java:

public static final ConfigOption<Boolean> TABLE_EXEC_HIVE_SINK_STATISTIC_AUTO_GATHER_ENABLE = key("table.exec.hive.sink.statistic-auto-gather.enable") .booleanType() .defaultValue(true) // 1.16 起默认为 true .withDescription( "If it's true, Flink will gather statistic automatically during " + "writing Hive Table. ... For ORC and Parquet format, " + "numFiles/totalSize/numRows/rawDataSize can be gathered. " + "For other format, only numFiles/totalSize can be gathered. " + "Note: only batch mode supports auto gather statistic, " + "stream mode doesn't support it yet.");

源码确认了发布说明未展开的几个细节:

  • 默认值确认为true,且仅批模式生效,流模式尚不支持自动收集;
  • 统计能力与存储格式相关:ORC/Parquet 可收集numFiles/totalSize/numRows/rawDataSize四项,其他格式只能收集numFiles/totalSize;
  • 还配套了一个并发选项table.exec.hive.sink.statistic-auto-gather.thread-num,默认 3 个线程用于 ORC/Parquet 统计收集,写入分区很多时可调大。

升级建议:批作业若写出的文件数量庞大且下游不依赖 Metastore 统计(如靠 Spark 的 file 级统计做优化),可在flink-conf.yaml中显式关闭:

table.exec.hive.sink.statistic-auto-gather.enable: false

Hive 版本支持范围收窄

发布说明(FLINK-27044):Flink 不再支持 Hive 1.x、2.1.x 与 2.2.x,原因是这些版本已不被 Hive 社区维护。仓库中可以看到目前保留的 SQL connector 产物:flink-sql-connector-hive-2.3.10 与 flink-sql-connector-hive-3.1.3,即官方预打包的 Hive 版本只剩 2.3.10 与 3.1.3 两档。

升级影响:若作业使用 Hive 1.x/2.1/2.2 的 MetaStore/HCatalog,1.16 上需要自行升级 Hive 环境,或从源码构建对应 connector 并自行保证稳定性。

Elasticsearch connector 迁移至独立仓库

发布说明(FLINK-26884):Elasticsearch connector 已从 Flink 主仓库复制到独立的 connector 仓库独立版本化维护。1.16 发布周期内两个仓库的产物同时存在但版本号体系不同:随 Flink 主线的产物版本为1.16.0,外部独立维护的产物版本为3.0.0。官方建议开发者在本发布周期内迁移到后者,以对齐 Flink 官方 connector 独立仓库(如 Kafka、JDBC 等)的版本节奏。

注意事项:迁移时要同步更换依赖坐标与版本,且外部版本3.0.0面向 Flink 1.16 运行环境,升级前请核对依赖树中不再混用新旧两套 Elasticsearch connector 类。

Pulsar Connector:cursor API 破坏性变更

发布说明(FLINK-27399)列出了 Pulsar connector cursor API 的三处破坏性变更:

  • 移除CursorPosition#seekPosition();
  • 移除StartCursor#seekPosition();
  • StopCursor#shouldStop的返回值由boolean改为StopCondition。

升级影响:仅影响自行实现或扩展了 Pulsar cursor 接口的用户,直接使用 connector 官方实现的作业不受影响;自定义 cursor 需要按新接口签名重写。

StreamingFileSink 正式标记弃用

发布说明(FLINK-27188):StreamingFileSink被标记为 deprecated,取而代之的是自 Flink 1.12 起逐步统一的FileSink。若项目仍在使用StreamingFileSink(及其配套的状态化路径策略 API),建议升级到FileSink.forRow(...)/FileSink.forBulkFormat(...)的新 API,后者在文件滚动、状态恢复语义上已覆盖旧实现。

Avro 生成代码命名空间修正

发布说明(FLINK-25962):Flink 生成的 Avro schema 命名空间改为org.apache.flink.avro.generated,以兼容 Avro Python SDK(此前生成的 schema 命名空间在 Python 侧无法被正确解析)。使用 Avro format 且与 Python 生态互通的场景,升级后该问题将自动修复;Java 侧若对生成类全限定名有硬编码引用,需同步调整。

AsyncSink 支持可配置 RateLimitingStrategy

发布说明(FLINK-28487):AsyncSinkWriter新增可配置的RateLimitingStrategy,sink 实现方可以针对特定 sink 定制请求失败时的限流退避行为;不指定时默认沿用AIMDRateLimitingStrategy(即 Additive-Increase/Multiplicative-Decrease,类似 TCP 拥塞控制思路的自适应限流)。

这条变更主要面向 sink 的开发者:基于 AsyncSink API 封装外部系统写入时,可以在失败率高的场景下选择更保守或更激进的速率策略,而无需 fork 运行时逻辑。

Runtime & Coordination

Metrics Reporter 按类名加载方式弃用

发布说明(FLINK-27206):通过metrics.reporter.<name>.class指定 reporter 类名的配置方式被弃用,reporter 实现应当提供MetricReporterFactory,所有配置迁移到 factory 机制下。特别地:若 reporter 从 plugins 目录加载,metrics.reporter.<name>.class将直接不再生效(因为插件类加载隔离下无法跨加载器按类名实例化)。

迁移检查:全局检索配置中的metrics.reporter.*.class,确认对应 reporter jar 已提供MetricReporterFactory(通过META-INF/services/org.apache.flink.metrics.MetricReporterFactorySPI 注册)。

Datadog reporter 的tags选项弃用

发布说明(FLINK-29002):DatadogReporter的tags选项被弃用,改用通用的scope.variables.additional选项(由 metrics core 的 scope 机制统一下发额外标签)。仓库中 Datadog reporter 位于 flink-metrics-datadog,可在其中确认 factory 化的配置入口。

REST API:未就绪返回 503 而非 500

发布说明(FLINK-25269):当 REST 请求到达但后端组件尚未就绪时,1.16 返回503 Service Unavailable(此前返回 500)。这是一个面向客户端的语义修正:监控探针、CI 中"等待 Flink 就绪"的脚本应以 503 作为"还在启动"的正常信号做重试,而把 500 视为真正异常。

Application 模式 + HA 时 JobID 改为基于 cluster ID 生成

发布说明(FLINK-19358):开启 HA 的 application mode 下,JobID 不再是全零的0000000000...,而是基于 cluster ID 生成。

升级影响:此前依赖"application mode 的 JobID 恒为 000..."这一隐含约定的外部工具(如脚本中写死 JobID、按 JobID 匹配告警/状态路径的逻辑)需要改为通过 REST API 动态查询 JobID。

Checkpoints

引入 Overdraft Buffer 缓解非对齐 Checkpoint 阻塞

发布说明(FLINK-26762)是 1.16 网络栈层面最重要的改进之一:为缓解反压期间 subtask 线程被"不可中断地阻塞"、导致 unaligned checkpoint barrier 无法注入的问题,1.16 引入了超支网络缓冲区(overdraft buffers)概念——subtask 可以在正常配置的 buffer 数量之外,额外向 BufferPool 申请默认最多 5 个超支 buffer,保证在背压高峰时仍有空间接收 checkpoint barrier 数据。

该行为的开关与额度定义在 NettyShuffleEnvironmentOptions.java:

taskmanager.network.memory.max-overdraft-buffers-per-gate # 默认 5,设为 0 可恢复 1.15 旧行为

升级建议:

  1. 该变更会略微增加作业内存占用(每个 gate 最多 5 个额外 buffer),大规模网络 buffer 配置的作业建议在压测环境核对 TaskManager 峰值内存;
  2. 若升级后出现 OOM 或内存压力,可显式设置taskmanager.network.memory.max-overdraft-buffers-per-gate: 0恢复 1.15 前的行为,代价是 unaligned checkpoint 在强反压下重新暴露阻塞风险;
  3. 该特性与 unaligned checkpoint 配置(execution.checkpointing.unaligned.*)配合使用效果最佳,二者共同构成 1.16 反压场景下 checkpoint 可用性的改进闭环。

table.exec.uid.generation:修复 1.15.x 非确定性 UID 问题

发布说明(FLINK-28861)揭示了一个影响面较大的兼容性陷阱:Flink 1.15.0 与 1.15.1 为从非编译计划(non-compiled plans)构建的算子生成了非确定性 UID,导致同一作业两次部署的算子 UID 不一致,从而使基于 UID 的状态恢复/版本升级变得困难甚至不可能。

1.16 通过新配置项table.exec.uid.generation修复,其默认行为是不再为这类新管线设置 UID;如果用户在 1.15.0/1 上已经接受了当时的 UID 行为并需要与之对齐,可显式设置:

table.exec.uid.generation: ALWAYS

升级路径建议:

  • 若当前跑在 1.15.0/1.15.1 且依赖其 UID 做状态恢复,升级 1.16 前务必先确认 Savepoint 中算子 UID 与新版本的匹配情况,必要时设置ALWAYS保持兼容;
  • 1.16 之后新建的 SQL 作业使用默认行为即可,避免 UID 在不同部署间漂移。

Python:PyFlink 1.16 将是最后支持 Python 3.6 的版本

发布说明(FLINK-28195):Python 3.6 的扩展支持已于 2021 年 12 月 23 日结束,官方计划PyFlink 1.16 为最后一个支持 Python 3.6 的版本。

升级行动项:仍使用 Python 3.6 的 PyFlink 环境,在升级至 1.17 及以后版本前必须先完成 Python 运行时到 3.7+(推荐 3.8/3.9)的迁移。仓库中 flink-python 模块的依赖声明可以辅助核对各版本兼容矩阵。

依赖升级

Hadoop 文件系统实现升级至 3.3.2

发布说明(FLINK-27308):Flink 文件系统实现所依赖的 Hadoop 实现升级至 3.3.2,带来两方面收益:

  • 获得 Hadoop 3.3.x 对应的文件系统特性;
  • 支持客户端侧的 Flink 状态加密(HADOOP-13887 对应能力),即 KMS 加密的 FS 上读写 Flink 状态时无需依赖服务端透明处理。

使用 S3/GCS/Azure/OSS 等通过 hadoop-fs 实现接入的集群,升级后可关注对应 filesystem 模块(见 flink-filesystems 目录下的各 hadoop 实现子模块)带来的行为变化。

Kafka Client 升级至 3.1.1

发布说明(FLINK-28060):Kafka connector 默认使用Kafka client 3.1.1。

注意事项:Kafka client 3.x 要求 broker 端至少为 0.10.0;若集群使用较老的 broker,或依赖了被旧 client 容忍的特殊行为,升级前应在测试环境验证 consumer 组、offset 提交与序列化路径。

Hive 2.3 connector 升级至 2.3.9

发布说明(FLINK-27063):Hive 2.3 connector 版本从旧版升级至2.3.9,与前面提到的"保留 Hive 2.3.10/3.1.3 预打包 connector"共同构成 1.16 的 Hive 支持矩阵。

PyFlink 依赖版本更新(支持 Python 3.9 与 Apple M1)

发布说明(FLINK-25188):为支持 Python 3.9 与 M1(ARM64)架构,PyFlink 更新了一组依赖:

apache-beam==2.38.0 arrow==5.0.0 pemja==0.2.6

自建 PyFlink wheel 或维护内部镜像的团队,同步升级这三项依赖即可对齐官方行为(尤其pemja是 PyFlink 与 JVM 之间对象桥接的核心组件,版本需与 Flink 发行版严格匹配)。

系统资源 Metrics 依赖更新

发布说明同时列出系统资源指标相关依赖的升级:

com.github.oshi:oshi-core:6.1.5 (MIT License) net.java.dev.jna:jna-platform:5.10.0 net.java.dev.jna:jna:5.10.0

影响面为cpu.load、内存等系统级 metric 的采集组件;常规部署无需操作,仅在自打包发行版(shaded dist)时需要更新依赖树。

总结:1.16 升级检查清单

结合上述各节,按"必须先做 / 需要评估 / 可选优化"三档整理升级检查清单:

类别事项依据条目
必须先做移除 String Expression DSL 用法,改用类型安全 Expression APIFLINK-26704
必须先做Hive 1.x/2.1.x/2.2.x 用户升级 Hive 环境或改用 2.3.10/3.1.3 connectorFLINK-27044
必须先做Python 3.6 用户规划迁移至 3.7+(1.17 起不再支持)FLINK-28195
必须先做重写自定义 Pulsar cursor(seekPosition已移除)FLINK-27399
需要评估Hive 批 Sink 统计自动收集是否关闭(statistic-auto-gather.enable)FLINK-28883
需要评估1.15.0/1 升级路径的 UID 匹配(table.exec.uid.generation)FLINK-28861
需要评估强反压作业内存余量 vs. overdraft buffer(默认 5,可设 0 回退)FLINK-26762
需要评估metrics reporter 从.class配置迁移到MetricReporterFactoryFLINK-27206
需要评估REST 客户端将 503 处理为"未就绪"而非错误FLINK-25269
可选优化Append-only 维表查询作业开启table.exec.async-lookup.output-mode=ALLOW_UNORDEREDFLINK-27622
可选优化维表查询不稳定场景启用 Retryable Lookup Join HINTFLINK-28779
可选优化jobmanager.sh的 host/web-ui-port 参数迁移为-D动态属性FLINK-28735

Flink 1.16 的整体基调是"平滑但有牙齿":绝大多数变更是新增能力与加固,但 String DSL 移除、Pulsar cursor 破坏性变更、Hive 版本收窄和 PyFlink 3.6 的退役四个点是硬断点,升级窗口内应优先安排上述检查项的验证。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:AirLLM 深度拆解:4GB 显存下 70B 大模型推理的最小启动配置
下一篇:终极Emissary-Ingress金丝雀部署教程:如何实现零停机应用发布

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询