ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink 1.7 升级指南:从 1.6 迁移到 1.7 的关键行为变更、配置调整与兼容性说明

Flink 1.7 升级指南:从 1.6 迁移到 1.7 的关键行为变更、配置调整与兼容性说明 大数据流处理批处理数据工程【免费下载链接】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 引入的 SAMSingle 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 evolutionTypeSerializerSnapshot全面取代旧序列化器快照机制Flink 1.7 之前序列化器快照通过TypeSerializerConfigSnapshot实现且序列化器 schema 兼容性检查逻辑内嵌在TypeSerializer的ensureCompatibility(TypeSerializerConfigSnapshot)方法中。1.7 引入了新的TypeSerializerSnapshot接口旧的TypeSerializerConfigSnapshot已被标记为废弃deprecated并将在未来版本中彻底移除。在当前仓库源码中TypeSerializerSnapshotT接口标注为PublicEvolving定义了如下核心方法int getCurrentVersion()返回当前快照二进制格式的版本号void writeSnapshot(DataOutputView out)将序列化器配置快照写入输出流void readSnapshot(int readVersion, DataInputView in, ClassLoader userCodeClassLoader)按写入版本读取快照支持跨版本的格式演进TypeSerializerT 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。在仓库源码中查询服务通过RpcMetricQueryServiceRetrieverflink-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 方式启动 TaskManager1.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。SQLOVER 窗口 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.xflink-conf.yamlSavepoint 管理恢复期间不移动/删除最近 savepoint运维流程指标查询端口开放metrics.internal.query-service.portflink-conf.yaml / 防火墙延迟指标显式设置metrics.latency.interval与metrics.latency.granularityflink-conf.yaml 或ExecutionConfigHadoop Netty不再假设io.netty在 shaded jar 中作业依赖声明本地恢复建议启用state.backend.local-recovery: trueflink-conf.yamlHA job-mode每个作业/集群设置独立high-availability.cluster-idflink-conf.yamlfailover 策略移除或设为full有状态作业flink-conf.yamlOVER 窗口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点击查看免费下载相关推荐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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进