Flink 1.7 升级指南:从 1.6 迁移到 1.7 的关键行为变更、配置调整与兼容性说明
2026/9/23 7:22:10 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

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

本指南基于当前仓库中的 Flink 1.7 版本发布说明(docs/content.zh/release-notes/flink-1.7.md),系统梳理从 Flink 1.6 升级到 1.7 时涉及的核心变更:Scala 2.12 编译兼容、状态序列化框架演进(TypeSerializerSnapshot)、savepoint 恢复语义、指标系统配置、本地恢复与多 slot TaskManager 支持等。阅读完本文,你将掌握升级 1.7 前必须完成的代码调整、flink-conf.yaml配置项核对清单,以及若干已知限制(如 Scala 2.12 下 shell 不可用、非默认 failover 策略的局限)的规避方案。

适用前提:本文内容以当前仓库(Flink 1.7 分支代码)为准,主要面向计划从 Flink 1.6.x 升级到 1.7 的集群运维与作业开发人员。

Scala 2.12 支持:lambda 推断变化带来的显式类型标注要求

Flink 1.7 开始提供 Scala 2.12 构建。由于 Scala 2.12 改变了 lambda 的实现方式——它现在使用 Java 8 引入的 SAM(Single Abstract Method,单抽象方法)接口支持——导致部分方法调用在同时存在 Scala 风格 lambda 与 SAM 候选时产生歧义。因此,在 Scala 2.12 下,一些原本无需显式类型标注的位置现在必须补充。

仓库中的TransitiveClosureNaive示例(flink-examples/flink-examples-batch/src/main/java/org/apache/flink/examples/java/graph/TransitiveClosureNaive.java)对应的 Scala 版本展示了这一变化。升级前(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 编译,请在升级后重新编译全部 Scala API 代码,重点检查coGroupjoincross等高阶函数调用处,编译器会明确提示需要补充类型标注的位置。

State evolution:TypeSerializerSnapshot全面取代旧序列化器快照机制

Flink 1.7 之前,序列化器快照通过TypeSerializerConfigSnapshot实现,且序列化器 schema 兼容性检查逻辑内嵌在TypeSerializerensureCompatibility(TypeSerializerConfigSnapshot)方法中。1.7 引入了新的TypeSerializerSnapshot接口,旧的TypeSerializerConfigSnapshot已被标记为废弃(deprecated),并将在未来版本中彻底移除。

在当前仓库源码中,TypeSerializerSnapshot<T>接口(标注为@PublicEvolving)定义了如下核心方法:

  • int getCurrentVersion():返回当前快照二进制格式的版本号;
  • void writeSnapshot(DataOutputView out):将序列化器配置快照写入输出流;
  • void readSnapshot(int readVersion, DataInputView in, ClassLoader userCodeClassLoader):按写入版本读取快照,支持跨版本的格式演进;
  • TypeSerializer<T> restoreSerializer():根据快照重建序列化器实例,用于安全读取旧数据;
  • resolveSchemaCompatibility(...):检查新序列化器读取旧数据格式时的兼容性,结果可以是完全兼容、需要重新配置(reconfigure)、格式不兼容,或需要迁移(migration,即用快照产生的序列化器反序列化旧数据、再用新序列化器重写)。

新的设计把"序列化器的二进制格式 schema"与"序列化器实例"解耦,使状态序列化与 schema 演进具备面向未来的灵活性。官方强烈建议从旧抽象迁移到新接口(具体迁移指南见上游文档 custom_serialization 章节),这样才能在未来版本中平滑演进状态序列化器与状态 schema。

实践要点:升级 1.7 后,凡是自定义了TypeSerializer的作业,应同步实现新的TypeSerializerSnapshot,否则在状态恢复时可能触发兼容性告警或失败。仓库中还提供快照读写工具类用于版本化读写。

Legacy mode 移除

Flink 1.7 不再支持 legacy 模式(legacy mode)。如果作业或集群配置仍依赖该模式,请继续使用 Flink 1.6.x,或在升级前移除相关配置。

Savepoint 参与故障恢复:语义与运维注意事项

1.7 之前,使用 exactly-once 语义的 sink 在"savepoint 完成后、下一次 checkpoint 完成前"发生故障时,可能出现重复输出数据。1.7 起,savepoint 会被用于恢复(recovery)流程,这意味着:

  • savepoint 不再完全由用户独占控制;
  • 如果之后没有更新的 checkpoint 或 savepoint,则不应移动或删除现有的 savepoint,否则会导致恢复时找不到可用的恢复点。

运维建议:建立严格的 savepoint 生命周期管理,确保恢复点目录在作业运行期间保持稳定。

MetricQueryService 独立线程池:新增端口配置

1.7 中,metric query service 运行在自己的ActorSystem中,因此需要开放一个新的端口供各 query service 之间通信。该端口通过flink-conf.yaml中的metrics.internal.query-service.port配置(对应文档锚点#metrics-internal-query-service-port)。在仓库源码中,查询服务通过RpcMetricQueryServiceRetriever(flink-runtime/src/main/java/org/apache/flink/runtime/webmonitor/retriever/impl/RpcMetricQueryServiceRetriever.java)在集群入口处被创建并传递(见 ClusterEntrypoint.java 中metricRegistry.getMetricQueryServiceRpcService()的调用)。

# flink-conf.yaml metrics.internal.query-service.port: <port>

如果端口未开放,JobManager 与 TaskManager 的指标查询服务将无法互通,Web 界面与 REST API 上的指标读取可能异常。

延迟指标粒度默认值变更

1.7 修改了延迟指标(latency metrics)的默认粒度。若要恢复 1.6 的行为,必须在flink-conf.yaml中显式将metrics.latency.granularity设置为subtask(对应文档锚点#metrics-latency-granularity)。

延迟标记默认关闭

延迟指标(latency metrics)在 1.7 中默认禁用。所有未通过ExecutionConfig#setLatencyTrackingInterval显式设置延迟追踪间隔的作业都将受到影响(不再上报延迟指标)。在仓库源码 ExecutionConfig.java 中,setLatencyTrackingInterval(long interval)用于设置该间隔(单位毫秒),且可在 Flink 配置中通过metrics.latency.interval覆盖。

若需恢复 1.6 的默认行为,请在flink-conf.yaml中配置:

# flink-conf.yaml metrics.latency.interval: <interval_milliseconds> metrics.latency.granularity: subtask

或在代码中显式调用:

env.getConfig().setLatencyTrackingInterval(5000L); // 例如 5 秒

Hadoop Netty 依赖重定位

1.7 将 Hadoop 的 Netty 依赖从io.netty重定位到org.apache.flink.hadoop.shaded.io.netty。影响如下:

  • 你可以在作业中捆绑自己的 Netty 版本,不再与flink-shaded-hadoop2-uber-*.jar中的 Netty 冲突;
  • 不能再假设io.netty存在于flink-shaded-hadoop2-uber-*.jar中,代码若直接引用io.netty包下的类,需要显式声明 Netty 依赖。

本地恢复(Local Recovery)修复

得益于调度器的改进,启用本地恢复后,故障恢复不再需要比故障前更多的 slot。官方鼓励用户在flink-conf.yaml中启用本地恢复:

# flink-conf.yaml state.backend.local-recovery: true

本地恢复的目录语义在 working_directory.md 中有说明:启用后 TaskManager 会在本地目录保存状态副本,从而在单节点故障时加速恢复。

多 slot TaskManager 支持

Flink 1.7 开始正式支持带多个 slot 的 TaskManager。此前建议以单 slot 方式启动 TaskManager,1.7 起不再有此限制,可以按资源情况为每个 TaskManager 配置任意数量的 slot(通过taskmanager.numberOfTaskSlots配置)。

StandaloneJobClusterEntrypoint 使用固定 JobID

standalone-job.sh脚本启动、并用于 job-mode 容器镜像的StandaloneJobClusterEntrypoint,在 1.7 中以固定 JobID启动所有作业。脚本位于 flink-dist/src/main/flink-bin/bin/standalone-job.sh,其入口点为standalonejob

这一变更的直接后果是:若要以 HA 模式运行多个作业/集群,必须为每个作业/集群设置不同的high-availability.cluster-id(对应文档锚点#high-availability-cluster-id),否则多个作业共享同一个 JobID 会在 HA 元数据上相互干扰。ZooKeeper HA 场景下的配置示例见 zookeeper_ha.md:

# flink-conf.yaml high-availability: zookeeper high-availability.cluster-id: /cluster_one # 重要:每个集群必须自定义

已知限制一:Scala 2.12 下 Scala shell 不可用

Flink 的 Scala shell 在 Scala 2.12 下无法工作(详见上游 issue FLINK-10911)。因此,flink-scala-shell模块不会为 Scala 2.12 发布。使用 Scala 2.12 的用户应避免依赖该模块。

已知限制二:非默认 failover 策略的局限

Flink 1.7 的非默认 failover 策略仍是高度实验性的功能,附带一系列限制:

  • 仅适用于无状态流作业
  • 其他任何场景下,强烈建议从flink-conf.yaml中移除jobmanager.execution.failover-strategy配置项,或将其显式设为"full"
# flink-conf.yaml jobmanager.execution.failover-strategy: full

为避免用户踩坑,该功能已从 1.7 文档中移除,直到其被修复(详见上游 issue FLINK-10880)。故障恢复策略的完整背景可参考 task_failure_recovery.md。

SQL:OVER 窗口 preceding 子句变为可选

Flink 1.7 起,OVER 窗口的preceding子句变为可选,未指定时默认值为UNBOUNDED(无界)。这简化了累加窗口类 SQL 的书写,例如:

-- 1.7 之前需要显式指定范围 SELECT SUM(amount) OVER (ORDER BY ts ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) ... -- 1.7 起可省略 preceding(默认 UNBOUNDED) SELECT SUM(amount) OVER (ORDER BY ts) ...

OperatorSnapshotUtil 输出 v2 格式快照

1.7 中,使用OperatorSnapshotUtil创建的快照会以 savepoint 格式v2写入。升级后,通过该工具生成的外部恢复快照格式与旧版本不兼容,相关工具链(如依赖 v1 格式的脚本)需要同步升级。

SBT 项目与 MiniClusterResource:需要显式 test-jar 依赖

MiniClusterResource已从flink-test-utils迁移到flink-runtimeflink-test-utils本身对flink-runtime声明了test-jar依赖,但sbt 无法正确拉取传递的 test-jar 依赖(详见上游 sbt issue #2964)。因此,使用 sbt 的项目必须显式声明:

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

这样测试代码才能正确引用MiniClusterResource相关的类。

升级核对清单

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

检查项处理方式涉及配置/代码
Scala 2.12 编译为歧义 lambda 补充显式类型标注Scala API 代码
自定义序列化器迁移到TypeSerializerSnapshot新接口状态序列化代码
Legacy mode移除相关配置(否则停留 1.6.x)flink-conf.yaml
Savepoint 管理恢复期间不移动/删除最近 savepoint运维流程
指标查询端口开放metrics.internal.query-service.portflink-conf.yaml / 防火墙
延迟指标显式设置metrics.latency.intervalmetrics.latency.granularityflink-conf.yaml 或ExecutionConfig
Hadoop Netty不再假设io.netty在 shaded jar 中作业依赖声明
本地恢复建议启用state.backend.local-recovery: trueflink-conf.yaml
HA + job-mode每个作业/集群设置独立high-availability.cluster-idflink-conf.yaml
failover 策略移除或设为"full"(有状态作业)flink-conf.yaml
OVER 窗口preceding省略时默认UNBOUNDEDSQL 作业
快照工具确认兼容 v2 savepoint 格式外部工具链
SBT 测试显式添加flink-runtimetest-jar 依赖build.sbt

以上配置项均可在flink-conf.yaml中调整;具体参数说明与完整取值,可在仓库的部署配置文档中进一步查阅(如 docs/content.zh/docs/deployment/config.md 及 docs/content.zh/docs/ops/metrics.md)。

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

【免费下载链接】flink

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

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

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

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

立即咨询