
后端大数据流处理批处理【免费下载链接】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 引入的 SAMSingle 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 TypeSerializerSnapshotT { int getCurrentVersion(); void writeSnapshot(DataOuputView out) throws IOException; void readSnapshot(int readVersion, DataInputView in, ClassLoader userCodeClassLoader) throws IOException; TypeSerializerSchemaCompatibilityT resolveSchemaCompatibility(TypeSerializerSnapshotT oldSerializerSnapshot); TypeSerializerT 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 指标粒度调整默认值不再是 subtask1.7 修改了 latency 指标的默认粒度。若想恢复 1.6 的行为必须显式把metrics.latency.granularity设置为subtask。当前仓库中该选项定义于 MetricOptions.java可取值及语义为取值语义single不区分 source 与 subtask仅跟踪整体延迟operator区分 source但不区分 subtask当前默认值subtask同时区分 source 与 subtask1.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。这带来两个实际影响你可以在自己的作业中打入任意版本的 Netty不必再担心与flink-shaded-hadoop2-uber-*.jar中的 Netty 冲突不能再假设flink-shaded-hadoop2-uber-*.jar中存在io.netty——依赖该包内 Netty 的代码需要改用重定位后的包路径或显式声明自己的 Netty 依赖。十、Local recovery 修复调度改进后恢复不再需要更多 slot1.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 backendEmbeddedRocksDBStateBackend与HashMapStateBackend本地状态根目录由execution.checkpointing.local-backup.dirs历史键名taskmanager.state.local.root-dirs指定默认落在WORKING_DIR/localState。十一、多 slot TaskManager 得到完整支持1.7 正式完整支持多 slot 的 TaskManagerTaskManager 可以以任意数量的 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.javahigh-availability.cluster-id默认值/default用于在 HA 存储中区分多个 Flink 集群standalone 集群需要显式设置YARN 模式下会自动推断。配置示例high-availability.cluster-id: /my-job-cluster十三、已知限制Scala shell 不兼容 Scala 2.12Flink 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 配置符合上述约束。十五、SQLOver 窗口的 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 建议按如下清单逐项核对Scala 2.12 用户检查coGroup/join等 lambda 写法补全显式类型标注自定义序列化器将TypeSerializerConfigSnapshot迁移到TypeSerializerSnapshot参照 custom_serialization.md移除 legacy 模式相关配置savepoint 生命周期管理不再随意移动/删除旧 savepoint确认metrics.internal.query.service.port端口段已开放、集群内不冲突按需设置metrics.latency.granularity与metrics.latency.interval以恢复延迟指标检查作业是否依赖flink-shaded-hadoop2-uber-*.jar中的io.netty类路径改用重定位包或自带 Netty评估并启用state.backend.local-recovery多 slot TaskManager 部署可放开单 slot 限制使用standalone-job.sh HA 时为每个 job/cluster 设置独立的high-availability.cluster-id有状态作业的 failover 策略保持full或删除相关配置SBT 项目补充flink-runtimetest-jar 依赖。全部变更的官方出处为 docs/content/release-notes/flink-1.7.md涉及的可配置项定义与源码证据已在上文逐一标注可据此深入仓库核对实现细节。赞分享后端大数据流处理批处理【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink 1.7 升级指南从 1.6 迁移到 1.7 的关键行为变更、配置调整与兼容性说明Flink 1.7 升级指南从 1.6 迁移到 1.7 的关键行为变更、配置调整与兼容性说明 本指南基于当前仓库中的 Flink 1.7 版本发布说明 do后端大数据流处理批处理Apache Flink 1.14 升级指南从 1.13 迁移的关键变更、配置与行为详解Apache Flink 1.14 升级指南从 1.13 迁移的关键变更、配置与行为详解 本指南基于 Flink 1.14 官方 Release Notes后端大数据流处理批处理终极TinyColor升级指南从1.5到1.6版本的关键变更与迁移策略终极TinyColor升级指南从1.5到1.6版本的关键变更与迁移策略 TinyColor是一款轻量级的JavaScript颜色处理库专注于提供快速、高效的前端开发工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考