☰
Flink 1.7 升级指南:从 1.6 迁移的关键变更、配置调整与源码级解读
2026/9/26 1:49:34 网站建设 项目流程
  • 后端
  • 大数据
  • 流处理
  • 批处理

【免费下载链接】flink

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

Flink 1.7 是一次面向稳定性和状态管理的重要版本升级:它引入了全新的TypeSerializerSnapshot状态序列化抽象、将 savepoint 纳入恢复流程、修复了本地恢复(local recovery)调度问题,同时移除了 legacy 模式。本文以官方 Release Notes 为骨架,逐一解析 Flink 1.6 升级到 1.7 时必须关注的行为变更、配置参数与依赖调整,并结合当前开源仓库中的源码(如 MetricOptions.java、TypeSerializerSnapshot.java)验证底层实现,帮助你平滑完成版本迁移。

一、升级前必读:这份 Release Notes 的定位

官方对这份文档的定位非常明确:它讨论的是 Flink 1.6 与 Flink 1.7 之间在配置、行为、依赖三个维度上的重要差异。如果你正计划将 Flink 版本升级到 1.7,务必逐条核对以下变更——其中既有会破坏编译的 Scala API 调整,也有会改变集群运行时语义的 savepoint 与指标行为,还有需要显式声明的新依赖。

仓库中该文档位于 docs/content/release-notes/flink-1.7.md,与 1.5、1.6、1.8 直至 1.20 的发布说明同目录存放,格式与篇幅保持一致。

二、Scala 2.12 支持:lambda 实现变化带来的显式类型标注

Flink 1.7 正式支持 Scala 2.12,但升级过程中可能需要在此前不需要标注类型的地方补上显式类型注解。官方以 Flink 代码库中的TransitiveClosureNaive.scala示例(当前仓库对应 Java 版本见 TransitiveClosureNaive.java)说明这种变化。

Scala 2.11 下的原代码:

val terminate = prevPaths .coGroup(nextPaths) .where(0).equalTo(0) { (prev, next, out: Collector[(Long, Long)]) => { val prevPaths = prev.toSet for (n <- next) if (!prevPaths.contains(n)) out.collect(n) } }

Scala 2.12 下必须改为:

val terminate = prevPaths .coGroup(nextPaths) .where(0).equalTo(0) { (prev: Iterator[(Long, Long)], next: Iterator[(Long, Long)], out: Collector[(Long, Long)]) => { val prevPaths = prev.toSet for (n <- next) if (!prevPaths.contains(n)) out.collect(n) } }

原因:Scala 2.12 改变了 lambda 的实现方式——现在利用 Java 8 引入的 SAM(Single Abstract Method)接口来实现 lambda。这导致一部分方法调用变得有歧义:原先只有 Scala 风格 lambda 是候选,现在 Scala lambda 和 SAM 同时成为候选方法,编译器无法再像以前那样唯一确定应调用哪个方法。

升级建议:凡是coGroup、join等接收函数式接口的高阶方法,在迁移到 Scala 2.12 后若出现"ambiguous reference to overloaded definition"类编译错误,即为prev、next这类参数补上Iterator[(Long, Long)]显式类型标注即可解决。

三、State Evolution:用 TypeSerializerSnapshot 全面取代 ConfigSnapshot

这是 1.7 版本在状态管理层面最重要的架构升级。

3.1 旧抽象为何被淘汰

在 Flink 1.7 之前,序列化器快照以TypeSerializerConfigSnapshot形式实现(该类型现已被标注@Deprecated,并将在未来版本中完全移除,由 1.7 引入的TypeSerializerSnapshot接口全面取代)。同时,序列化器 schema 兼容性检查的职责落在TypeSerializer自身,通过TypeSerializer#ensureCompatibility(TypeSerializerConfigSnapshot)方法实现。

旧抽象的问题在于:兼容性逻辑与序列化器强耦合,无法支撑 schema 的长期演进与序列化器的平滑迁移。

3.2 新抽象:TypeSerializerSnapshot

新的TypeSerializerSnapshot接口定义在 TypeSerializerSnapshot.java,其核心契约(官方文档 custom_serialization.md 有完整描述)为:

public interface TypeSerializerSnapshot<T> { int getCurrentVersion(); void writeSnapshot(DataOuputView out) throws IOException; void readSnapshot(int readVersion, DataInputView in, ClassLoader userCodeClassLoader) throws IOException; TypeSerializerSchemaCompatibility<T> resolveSchemaCompatibility(TypeSerializerSnapshot<T> oldSerializerSnapshot); TypeSerializer<T> restoreSerializer(); }

对应的TypeSerializer侧提供snapshotConfiguration()方法返回快照。新抽象将"写入快照时的序列化 schema"与"恢复时的兼容性判定"解耦:

  • getCurrentVersion/writeSnapshot/readSnapshot:管理快照自身的版本化读写,快照的写入格式演进不再受序列化器制约;
  • resolveSchemaCompatibility(oldSerializerSnapshot):在恢复时把旧序列化器快照交给新序列化器快照做兼容性判定,返回三种结果之一——compatibleAsIs()(schema 一致,可直接复用)、compatibleAfterMigration()(schema 不同但可用旧序列化器读、新序列化器写完成迁移)、incompatible()(无法迁移);
  • restoreSerializer():作为工厂,在需要迁移时重建出能识别旧 schema 的序列化器实例。

当前仓库中的序列化测试(如 TypeSerializerSnapshotTest.java)以及 CompositeTypeSerializerSnapshot.java(其中保留了多处@Deprecated的旧接口桥接代码)都能印证新旧抽象并存期间的演进痕迹。

3.3 迁移建议

官方在 1.7 发布说明中明确指出:为了让状态序列化器与 schema 具备面向未来的演进能力,强烈建议从旧抽象迁移到新抽象,完整的迁移指南见仓库内文档 custom_serialization.md。在迁移完成前,旧代码仍可通过兼容层运行,但TypeSerializerConfigSnapshot已被标记废弃,不应再作为新代码的基类。

四、移除 legacy 模式

Flink 1.7 不再支持 legacy 模式。如果业务强依赖该模式,官方建议停留在 Flink 1.6.x。升级前请检查flink-conf.yaml中是否有与 legacy 模式相关的开关或配置,并将其清理。注意这里的 "legacy mode" 与 1.18 之后 hybrid shuffle 的 "new/legacy mode"(见 batch_shuffle.md)是两个不同的概念,勿混淆。

五、Savepoint 纳入恢复流程:所有权语义变化

Flink 1.7 起,savepoint 会在恢复过程中被使用。此前,使用 exactly-once 语义的 sink 时,若在 savepoint 之后、下一个 checkpoint 之前发生故障,可能出现重复输出数据的问题。1.7 通过让 savepoint 参与恢复流程修复了该问题,但由此带来一个重要的行为变化:

savepoint 不再完全由用户掌控。如果没有更新的 checkpoint 或 savepoint,则不应移动或删除已有的 savepoint。

这一语义在源码中也有对应体现:JobGraph中保存恢复设置(见 JobGraph.java 中的SavepointRestoreSettings字段),恢复时指定RestoreMode.NO_CLAIM等模式(见 SavepointRestoreSettings.java)。

运维建议:升级后对 savepoint 目录的清理与迁移操作要更加谨慎,建议为 savepoint 保留足够生命周期,避免在未产生新 checkpoint/savepoint 的情况下删除旧 savepoint 导致恢复失败。

六、MetricQueryService 运行在独立线程池并占用新端口

1.7 之前,metric query service(用于 Web UI 与 queryable state 拉取指标)与主 RPC 体系共用进程资源。1.7 起,metric query service 运行在独立的ActorSystem中,因此需要为各 query service 之间的通信开放新的端口。

对应的配置键为metrics.internal.query.service.port,在flink-conf.yaml中设置。当前仓库中该选项定义于 MetricOptions.java,要点如下:

  • 默认值为"0",表示由 Flink 自动寻找空闲端口;
  • 支持单端口(如50100)、端口段(如50100-50200)或二者组合("50100,50101");
  • 官方推荐配置一段端口范围,避免同一机器上多个 Flink 组件发生端口冲突。

示例:

metrics.internal.query.service.port: 50100-50200

同文件还定义了metrics.internal.query.service.thread-priority(默认 1,取值范围 1–10,注意增大该值可能拖垮主组件),供需要调整查询服务线程优先级的场景使用。

七、Latency 指标粒度调整:默认值不再是 subtask

1.7 修改了 latency 指标的默认粒度。若想恢复 1.6 的行为,必须显式把metrics.latency.granularity设置为subtask。

当前仓库中该选项定义于 MetricOptions.java,可取值及语义为:

取值语义
single不区分 source 与 subtask,仅跟踪整体延迟
operator区分 source,但不区分 subtask(当前默认值)
subtask同时区分 source 与 subtask(1.6 时代的默认行为)

配置示例(恢复 1.6 行为):

metrics.latency.granularity: subtask

八、Latency marker 默认关闭:延迟指标默认不再产生

与粒度调整配套,1.7 将latency 指标默认关闭:所有未显式通过ExecutionConfig#setLatencyTrackingInterval设置追踪间隔的作业,都不会再产生 latency 指标。要恢复此前的默认行为,需要在flink-conf.yaml中配置metrics.latency.interval。

当前仓库中该选项定义于 MetricOptions.java,默认值为0ms——设置为 0 或负数即禁用 latency 追踪,且官方注明开启该特性会显著影响集群性能。配置示例:

metrics.latency.interval: 5 s

另外,同文件中的metrics.latency.history-size(默认 128)控制每个算子保留的历史延迟测量条数,与上述两项共同决定 latency 指标的完整行为。

九、Hadoop 的 Netty 依赖重定位

1.7 对 Hadoop 的 Netty 依赖做了进一步重定位:由io.netty移入org.apache.flink.hadoop.shaded.io.netty。

这带来两个实际影响:

  1. 你可以在自己的作业中打入任意版本的 Netty,不必再担心与flink-shaded-hadoop2-uber-*.jar中的 Netty 冲突;
  2. 不能再假设flink-shaded-hadoop2-uber-*.jar中存在io.netty——依赖该包内 Netty 的代码需要改用重定位后的包路径,或显式声明自己的 Netty 依赖。

十、Local recovery 修复:调度改进后恢复不再需要更多 slot

1.7 修复了 local recovery 与调度器之间的联动问题。此前,开启本地恢复时,故障恢复可能需要比故障前更多的 slot(因为本地副本与远端状态可能位于不同节点)。随着调度逻辑的改进,这种情况不再发生。

官方在发布说明中明确鼓励用户在flink-conf.yaml中启用 local recovery,配置键为state.backend.local-recovery。当前仓库中该选项定义于 CheckpointingOptions.java(注:该键在当前版本中已标记@Deprecated,迁移到 StateRecoveryOptions.java 中的LOCAL_RECOVERY,历史键名仍兼容),默认false,启用示例:

state.backend.local-recovery: true

同时注意两点:本地恢复当前仅覆盖 keyed state backend(EmbeddedRocksDBStateBackend与HashMapStateBackend);本地状态根目录由execution.checkpointing.local-backup.dirs(历史键名taskmanager.state.local.root-dirs)指定,默认落在<WORKING_DIR>/localState。

十一、多 slot TaskManager 得到完整支持

1.7 正式完整支持多 slot 的 TaskManager:TaskManager 可以以任意数量的 slot 启动,官方不再建议只使用单 slot。这意味着此前单 slot 启动的规避性实践可以废弃,集群资源利用率有望提升。升级后可根据实际负载为每个 TaskManager 配置多个 slot(通过taskmanager.numberOfTaskSlots控制),并重新评估 slot 与并行度的配比。

十二、StandaloneJobClusterEntrypoint 生成固定 JobID

由脚本standalone-job.sh(当前仓库见 standalone-job.sh,用于 job 模式容器镜像)启动的StandaloneJobClusterEntrypoint,从 1.7 起会为作业生成固定的JobID。

这个行为对 HA 部署有直接影响:每个 job/cluster 必须设置不同的high-availability.cluster-id,否则多个 job 会因共享 JobID 与 HA 命名空间而发生冲突。

当前仓库中该选项定义于 HighAvailabilityOptions.java:high-availability.cluster-id,默认值/default,用于在 HA 存储中区分多个 Flink 集群;standalone 集群需要显式设置,YARN 模式下会自动推断。配置示例:

high-availability.cluster-id: /my-job-cluster

十三、已知限制:Scala shell 不兼容 Scala 2.12

Flink 1.7 的Scala shell 无法在 Scala 2.12 下工作,因此flink-scala-shell模块不会为 Scala 2.12 发布。使用 Scala shell 的用户需继续停留在 Scala 2.11 版本线,或改用 SQL Client / REPL 类替代方案。该问题在 Apache JIRA 中跟踪为 FLINK-10911,待其修复后方可解除此限制。

十四、Failover 策略的已知局限:非默认策略仍为实验特性

1.7 中,非默认的 failover 策略仍是高度实验性的功能,附带一组已知限制:

  • 只有在运行无状态流作业时才建议使用该特性;
  • 其他任何情况下,强烈建议从flink-conf.yaml中删除jobmanager.execution.failover-strategy配置项,或将其显式设置为"full";
  • 为避免后续踩坑,该特性在修复前已从官方文档中移除。

当前仓库中该选项定义于 JobManagerOptions.java,接受的取值为full(重启全部 task 恢复作业)与region(仅重启受故障影响的 pipelined region,当前默认值)。问题跟踪见 Apache JIRA FLINK-10880。升级到 1.7 后,如作业涉及有状态算子,请确认 failover 配置符合上述约束。

十五、SQL:Over 窗口的 preceding 子句变为可选

Flink 1.7 的 SQL 语法调整:Over 窗口的preceding子句现在变为可选,未指定时默认为UNBOUNDED。

例如以下两种写法在 1.7 中等价:

-- 显式声明 UNBOUNDED PRECEDING SELECT a, SUM(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) FROM t; -- 省略 preceding 子句,隐含 UNBOUNDED SELECT a, SUM(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN CURRENT ROW AND CURRENT ROW) FROM t;

这一语法放宽降低了窗口 SQL 的书写负担。当前仓库中关于窗口帧(window frame)的完整语法与缺省规则(如ORDER BY存在但缺省window_frame时默认RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)可参考 window-functions.md。

十六、OperatorSnapshotUtil 改为写入 v2 格式

使用OperatorSnapshotUtil创建的快照,从 1.7 起以savepoint 格式 v2写入。该工具用于测试场景下的算子状态快照读写(当前仓库实现见 OperatorSnapshotUtil.java,其writeStateHandle方法直接以MetadataV3Serializer写入各类型状态句柄)。涉及以旧工具生成的测试快照文件时,需按新格式重新生成,否则读取可能失败。

十七、SBT 项目必须显式声明 flink-runtime 的 test-jar 依赖

如果 SBT 项目使用了MiniClusterResource(用于本地启动 MiniCluster 做集成测试),升级 1.7 后需要显式添加flink-runtime的 test-jar 依赖:

libraryDependencies += "org.apache.flink" %% "flink-runtime" % flinkVersion % Test classifier "tests"

原因在于:MiniClusterResource已从flink-test-utils迁移到flink-runtime(当前仓库位置见 MiniClusterResource.java,实现基于 JUnit 的ExternalResource规则)。虽然flink-test-utils正确地声明了对flink-runtime的 test-jar 传递依赖,但SBT 不会正确拉取传递的 test-jar 依赖(对应 sbt 社区 issue #2964),因此必须显式声明。Maven 用户通常不受影响,因为 Maven 会解析传递的 test-jar 依赖。

十八、升级检查清单

综合以上变更,从 1.6 升级到 1.7 建议按如下清单逐项核对:

  1. Scala 2.12 用户:检查coGroup/join等 lambda 写法,补全显式类型标注;
  2. 自定义序列化器:将TypeSerializerConfigSnapshot迁移到TypeSerializerSnapshot,参照 custom_serialization.md;
  3. 移除 legacy 模式相关配置;
  4. savepoint 生命周期管理:不再随意移动/删除旧 savepoint;
  5. 确认metrics.internal.query.service.port端口段已开放、集群内不冲突;
  6. 按需设置metrics.latency.granularity与metrics.latency.interval以恢复延迟指标;
  7. 检查作业是否依赖flink-shaded-hadoop2-uber-*.jar中的io.netty类路径,改用重定位包或自带 Netty;
  8. 评估并启用state.backend.local-recovery;
  9. 多 slot TaskManager 部署可放开单 slot 限制;
  10. 使用standalone-job.sh+ HA 时,为每个 job/cluster 设置独立的high-availability.cluster-id;
  11. 有状态作业的 failover 策略保持full或删除相关配置;
  12. SBT 项目补充flink-runtimetest-jar 依赖。

全部变更的官方出处为 docs/content/release-notes/flink-1.7.md,涉及的可配置项定义与源码证据已在上文逐一标注,可据此深入仓库核对实现细节。

  • 后端
  • 大数据
  • 流处理
  • 批处理

【免费下载链接】flink

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

相关推荐

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

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

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

立即咨询