- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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 代码,重点检查coGroup、join、cross等高阶函数调用处,编译器会明确提示需要补充类型标注的位置。
State evolution:TypeSerializerSnapshot全面取代旧序列化器快照机制
Flink 1.7 之前,序列化器快照通过TypeSerializerConfigSnapshot实现,且序列化器 schema 兼容性检查逻辑内嵌在TypeSerializer的ensureCompatibility(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-runtime。flink-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.port | flink-conf.yaml / 防火墙 |
| 延迟指标 | 显式设置metrics.latency.interval与metrics.latency.granularity | flink-conf.yaml 或ExecutionConfig |
| Hadoop Netty | 不再假设io.netty在 shaded jar 中 | 作业依赖声明 |
| 本地恢复 | 建议启用state.backend.local-recovery: true | flink-conf.yaml |
| HA + job-mode | 每个作业/集群设置独立high-availability.cluster-id | flink-conf.yaml |
| failover 策略 | 移除或设为"full"(有状态作业) | flink-conf.yaml |
| OVER 窗口 | preceding省略时默认UNBOUNDED | SQL 作业 |
| 快照工具 | 确认兼容 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
相关推荐
Apache Flink 1.14 升级指南:从 1.13 迁移的关键变更、配置与行为详解
Apache Flink 1.14 升级指南:从 1.13 迁移的关键变更、配置与行为详解 本指南基于 Flink 1.14 官方 Release Notes(
大数据流处理批处理数据工程NVD3版本迁移手册:从1.7.x到1.8.6的关键变更
NVD3版本迁移手册:从1.7.x到1.8.6的关键变更 NVD3作为基于D3.js的可复用图表库,从1.7.x到1.8.6版本的迭代包含多项重要变更,涉及AP
数据可视化前端hls.js 迁移指南:从 0.x 到 1.7 的完整升级路径与破坏性变更详解
hls.js 迁移指南:从 0.x 到 1.7 的完整升级路径与破坏性变更详解 本指南以 hls.js 官方迁移文档 MIGRATING.md https://
音视频前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考