ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink 1.8 版本升级指南:状态清理、序列化兼容、内存与配置变更全解析

Flink 1.8 版本升级指南:状态清理、序列化兼容、内存与配置变更全解析 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文基于 Apache Flink 1.8 官方 Release Notes 整理而成系统梳理 Flink 1.7 → 1.8 之间涉及配置、行为与依赖的关键变化覆盖状态管理TTL 持续清理、Savepoint Schema 迁移与兼容性、Maven 依赖Hadoop 打包方式、TaskManager 网络绑定策略、Table API 大规模接口重构、Kafka 连接器行为变更以及内存管理等多个方面。无论你是正在规划升级到 Flink 1.8还是需要理解 1.8 时期 API 设计决策对后续版本的影响本文都能帮助你提前识别升级风险、规避兼容性陷阱并掌握新特性的正确用法。说明文中涉及的实现类均可在当前仓库中找到对应源码master 分支实现并标注了具体文件路径供深入研读部分配置项在后续版本中发生了演进文中会同时说明 1.8 时期的行为与当前仓库中的实现差异。一、State状态相关变更1.1 Keyed State TTL 的持续增量清理Flink 1.6 引入了 Keyed State 的 TTLtime-to-live能力FLINK-9510允许在访问状态条目时清理并使其不可访问同时在进行 savepoint/checkpoint 写入时也会触发清理。但早期的 TTL 清理是惰性的——只在状态被访问或快照时触发无法保证过期条目被及时移除。Flink 1.8 为此引入了过期条目的持续增量清理continuous incremental cleanup分别针对 RocksDB 状态后端FLINK-10471与堆内存状态后端FLINK-10473实现。这意味着根据 TTL 设置判定为过期的旧条目会持续地被后台清理而不再依赖访问时机从而有效控制状态体积的无界增长。从当前仓库源码可以看到 TTL 机制的核心实现位于 TtlStateFactory.java该类将普通 state 对象包装为带 TTL 逻辑的 state 对象负责在读写路径上注入过期判定与清理动作。结合源码注释该工厂将后端产生的 state 对象包装上 TTL 逻辑可以推断 1.8 的持续清理正是在这层包装中按配置策略周期性扫描并驱逐过期条目。1.2 恢复 Savepoint 时的 Schema 迁移支持Flink 1.7.0 首次为AvroSerializer增加了修改状态 Schema 的能力FLINK-10605。Flink 1.8.0 在此基础上推进了所有内置TypeSerializer向新序列化器快照抽象serializer snapshot abstraction的迁移该抽象在理论上允许在恢复时进行 Schema 迁移。到 1.8 为止Flink 自带的序列化器中以下类型已正式支持 Schema 迁移序列化器覆盖类型JIRA 编号PojoSerializerPOJO 类型字段增删/调整FLINK-11485JavaEnumSerializer枚举常量增删FLINK-11334AvroSerializerAvro 记录 Schema 演进1.7 引入FLINK-10605Kryo有限场景下支持FLINK-11323以上序列化器在当前仓库中均可找到对应实现PojoSerializer.java、EnumSerializer.java、AvroSerializer.java。Schema 迁移的能力基础是新的序列化器快照抽象其核心类为 CompositeTypeSerializerSnapshot.java它允许多个嵌套序列化器在恢复时逐一比对并迁移。1.3 Savepoint 兼容性Scala TraversableSerializer 升级限制由于TraversableSerializer在 1.8 中发生更新FLINK-11539包含 ScalaTraversableSerializer的 Flink 1.2 Savepoint 将不再与 Flink 1.8 兼容。官方给出的绕行方案是先将版本升级到 Flink 1.3 1.7 之间的某个版本再升级到 Flink 1.8。即采用两步走升级路径避免跨版本序列化器不兼容导致的恢复失败。1.4 RocksDB 版本升级与切换到 FRocksDB为支撑带 TTL 的持续状态清理Flink 需要 RocksDB 提供特定的底层能力因此 1.8 切换到 RocksDB 的自定义构建版本FRocksDBFLINK-10471。所用 FRocksDB 构建基于升级后的RocksDB 5.17.2。平台兼容性注意对于 Mac OS XRocksDB 5.17.2 仅支持OS X 10.13 及以上版本。在更老的 macOS 上使用 Flink 1.8 RocksDB 状态后端时需要优先升级操作系统。二、Maven 依赖变化2.1 Hadoop 库不再随发行包捆绑FLINK-11266Flink 1.8 起包含 Hadoop 的便捷二进制发行包不再发布。这意味着官方下载页面默认提供的flink-dist不再内置flink-shaded-hadoop2部署方式发生如下变化若部署依赖flink-dist中内置的flink-shaded-hadoop2必须手动从官方下载页的可选组件optional components部分下载预打包的 Hadoop jar并将其复制到 Flink 安装目录的/lib目录下或自行构建含 Hadoop 的发行包对flink-dist执行打包并激活include-hadoopMaven profile即可在构建期将 Hadoop 打入发行包。由于 Hadoop 不再默认包含在flink-dist中原先打包时使用的-DwithoutHadoop参数不再对构建产生任何影响该开关已失去意义。升级到 1.8 后请检查自己的构建脚本与部署流水线移除或修正相关参数。三、配置变更3.1 TaskManager 默认绑定策略变化FLINK-11716Flink 1.8 中TaskManager 默认改为绑定主机 IP 地址而非主机名hostname。这一行为由配置项taskmanager.network.bind-policy控制。升级到 1.8 后如果集群出现难以解释的连接问题可以在flink-conf.yaml中显式设置以恢复 1.8 之前的行为taskmanager.network.bind-policy: name从当前仓库源码 TaskManagerOptions.java 可以看到该配置项的完整定义键为taskmanager.network.bind-policy类型为 string默认值为ip1.8 时期引入的新默认值其语义是在未显式设置taskmanager.host时TaskManager 自动选择绑定地址的策略可选值包括name—— 使用主机名作为绑定地址1.8 之前的行为ip—— 使用主机 IP 地址作为绑定地址1.8 起的新默认值。如果你的环境依赖主机名解析如 DNS 配置不完整、多网卡绑定需要主机名路由等场景请在升级时显式保留name策略避免连接中断。四、Table API 变更Flink 1.8 是 Table API 走向API 与实现分离的关键版本包含大量接口调整与弃用升级时需重点核对代码中的 Table API 用法。4.1 弃用直接使用Table构造函数FLINK-11447此前开发者可以直接通过Table的构造函数来执行与lateral table表函数的 join。Flink 1.8 弃用该用法应改用table.joinLateral(...)—— 普通侧向连接table.leftOuterJoinLateral(...)—— 左外侧向连接。这一变更的目的是将Table类转化为接口使 API 在未来更易维护、更干净。当前仓库中Table已是一组接口方法Table.java 上可以看到joinLateral(Expression tableFunctionCall)、joinLateral(Expression, Expression joinPredicate)以及对应的leftOuterJoinLateral重载注释中给出了典型用法例如table.joinLateral(call(MySplitUDTF.class, $(c)).as(s)); table.joinLateral(split($c) as s, $a $s);这印证了 1.8 的 API 演进方向侧向 join 统一收敛到显式的joinLateral/leftOuterJoinLateral方法上。4.2 引入符合 RFC 4180 的新 CSV 格式描述符FLINK-9964Flink 1.8 引入了符合 RFC 4180 规范的新 CSV 格式描述符类名为org.apache.flink.table.descriptors.Csv。需要注意的是新描述符目前只能与 Kafka 连接器配合使用。旧的描述符更名为org.apache.flink.table.descriptors.OldCsv继续用于文件系统连接器file system connectors。4.3 弃用 TableEnvironment 上的静态构建方法FLINK-11445为实现 API 与具体实现的分离TableEnvironment.getTableEnvironment()静态方法被弃用应改用BatchTableEnvironment.create(...)StreamTableEnvironment.create(...)。4.4 Table API 的 Maven 模块拆分FLINK-11064原先依赖flink-table的用户需要更新依赖声明改为依赖flink-table-planner并根据使用语言Java 或 Scala补充对应的 API bridge 依赖之一flink-table-api-java-bridgeJavaflink-table-api-scala-bridgeScala。当前仓库中 flink-table 目录下即可看到flink-table-planner、flink-table-api-java-bridge、flink-table-api-scala-bridge等模块的完整布局与 1.8 拆分后的结构一致。4.5 外部 Catalog 表构建器变更FLINK-11522ExternalCatalogTable.builder()被弃用改为使用ExternalCatalogTableBuilder()。4.6 Table API 连接器 jar 命名变更FLINK-11026Kafka、Elasticsearch 6 的SQL 连接器 jar 命名方案发生改变Maven 坐标中不再带sql-jar限定符artifactId 前缀由flink改为flink-sql。例如Kafka 连接器从原来的命名变为形如flink-sql-connector-kafka...。当前仓库 flink-sql-connector-hive-2.3.10、flink-sql-connector-hive-3.1.3 等模块的命名正是这一规则沿用的结果。4.7 Null 字面量写法变更FLINK-11785Table API 中定义 Null 字面量需改用nullOf(type)原先的Null(type)写法已弃用。五、Connectors连接器变更5.1 新增可直接访问 ConsumerRecord 的 KafkaDeserializationSchemaFLINK-8354针对 Flink 的KafkaConsumerFlink 1.8 引入了新的KafkaDeserializationSchema它可以直接访问 Kafka 的ConsumerRecord从而拿到消息的原始元数据topic、partition、offset、timestamp、headers 等。该新接口取代了KeyedSerializationSchema的功能后者虽已弃用但在 1.8 中仍然可用。如果你在自定义反序列化逻辑中需要基于消息头、offset 等元数据做处理应从KeyedSerializationSchema迁移到新的KafkaDeserializationSchema。5.2 FlinkKafkaConsumer 会按 Topic 过滤恢复的分区FLINK-10342从 Flink 1.8.0 起FlinkKafkaConsumer在从状态恢复时始终会过滤掉那些与当前订阅 Topic 规范不再关联的已恢复分区。这一行为在 1.8 之前的版本中不存在。若想保留旧行为可在FlinkKafkaConsumer上调用配置方法disableFilterRestoredPartitionsWithSubscribedTopics()官方文档给出了一个非常直观的场景示例假设某个 Kafka Consumer 原本消费 TopicA你做了 savepoint然后将该 Consumer 的订阅改为 TopicB再从该 savepoint 重启作业。变更之前由于状态中记录了正在消费 Topic A恢复后 Consumer 会同时消费 TopicA和B变更之后恢复时会用配置的 Topic 列表过滤状态中保存的 Topic因此只消费 TopicB。该变更避免了因状态残留导致的幽灵订阅但如果你依赖旧的跨 Topic 恢复语义务必显式调用上述禁用方法。六、其他接口变更6.1 从 TypeSerializer 接口移除 canEqual()FLINK-9803canEqual()方法通常用于在类型继承层级中做正确的相等性判断。由于TypeSerializer实际上并不需要该特性Flink 1.8 将其从接口中移除。实现自定义序列化器的用户需同步删除相关方法。6.2 移除 CompositeSerializerSnapshot 工具类FLINK-11073CompositeSerializerSnapshot工具类被移除对于将序列化委托给多个嵌套序列化器的复合序列化器快照应改用CompositeTypeSerializerSnapshot。该替代类的当前实现位于 CompositeTypeSerializerSnapshot.java官方建议阅读其中的实现与使用说明来迁移自定义复合序列化器。七、内存管理Flink 1.8.0 及之前版本中TaskManager 的托管内存managed memory占比由taskmanager.memory.fraction控制默认值为 0.7。这里存在一个经典陷阱JVM 参数NewRatio的默认值是 2意味着老年代old generation只占堆内存的 2/3约 0.66。当托管内存占比0.7超过老年代占比0.66时托管内存的一部分会被分配到新生代进而引发OOM内存溢出。因此如果升级后运行在taskmanager.memory.fraction 0.7默认配置下遇到 OOM官方建议手动将该值调低例如调整为 0.6 或更低使其低于老年代可用比例。补充说明在后续 Flink 版本中内存模型已全面重构托管内存占比配置演进为taskmanager.memory.managed.fraction。当前仓库 TaskManagerOptions.java 中该配置的默认值已变为0.4且基于Total Flink Memory计算同时支持taskmanager.memory.managed.size直接指定大小。1.8 时代的taskmanager.memory.fraction问题在新版本中已不再适用但理解这一历史背景有助于排查存量作业升级时遇到的 OOM。八、升级到 Flink 1.8 的检查清单综合以上变更规划升级时建议逐项核对状态与序列化确认作业中是否有 Flink 1.2 生成的、含 ScalaTraversableSerializer的 Savepoint需两步升级自定义序列化器需适配新快照抽象并移除canEqual()依赖检查部署是否依赖flink-dist内置 Hadoop需手动放入/lib或启用include-hadoopprofile 构建-DwithoutHadoop参数已失效配置确认集群网络环境是否兼容 TaskManager 默认的 IP 绑定策略必要时设置taskmanager.network.bind-policy: nameTable API替换弃用的Table构造函数、getTableEnvironment()、ExternalCatalogTable.builder()、Null(type)更新 Maven 依赖与连接器 jar 坐标Kafka 连接器确认恢复分区过滤的新行为是否符合预期必要时调用disableFilterRestoredPartitionsWithSubscribedTopics()内存如遇 OOM调低taskmanager.memory.fraction默认 0.7至低于老年代占比。按此清单完成核对后可显著降低升级到 Flink 1.8 过程中的行为回归与运行时故障风险。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink 升级指南应用状态兼容与跨版本 Savepoint 迁移实战Flink 升级指南应用状态兼容与跨版本 Savepoint 迁移实战 本文是 Flink 运维体系中的升级实操指南聚焦两个核心场景如何在不丢失状态的前提大数据流处理批处理数据工程Flink 1.10 升级指南从 1.9 迁移的关键变更、新内存模型与 RocksDB 状态管理Flink 1.10 升级指南从 1.9 迁移的关键变更、新内存模型与 RocksDB 状态管理 本文基于当前仓库 Flink 1.10 Release No大数据流处理批处理数据工程Dgraph版本升级兼容性指南API变更与适配Dgraph版本升级兼容性指南API变更与适配 你是否在升级Dgraph时遭遇过API调用失败、数据结构不兼容等问题本文将系统梳理Dgraph 23.x到2数据库图数据库分布式数据库后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表