Flink 1.16 发布说明深度解读:重试 Lookup Join、异步输出模式与网络超支缓冲等关键变更
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Apache Flink 1.16 是一次以稳定性加固 易用性改进为主线的版本。本篇以官方发布说明docs/content/release-notes/flink-1.16.md为主体完整梳理从 1.15 升级到 1.16 时需要关注的配置、行为与依赖变更并逐条对照当前仓库中的源码定义帮助你判断哪些变更会直接影响现有作业的升级路径。读完后你将能够回答三个问题哪些配置项的默认行为在 1.16 发生了变化、哪些 API 被移除需要迁移、以及新增的 Retryable Lookup Join / 异步输出模式 / 超支缓冲Overdraft Buffer分别如何在实际作业中配置。升级总览先记住这些默认值变化Flink 1.16 发布说明覆盖六大类变更Clusters Deployment、Table API SQL、Connectors、Runtime Coordination、Checkpoints、Python外加依赖升级。其中需要特别注意的默认行为变化即不改任何配置、升级后行为也会变有以下几处建议在升级演练中重点回归Hive Sink 在批模式下默认向 Hive Metastore 上报统计信息FLINK-28883非对齐 CheckpointUnaligned Checkpoint在反压场景下每个 gate默认允许申请 5 个超支网络缓冲区FLINK-26762可能略微增加作业内存占用Application 模式 HA 开启时JobID 不再固定为0000000000...而是基于 cluster ID 生成FLINK-19358REST API 在组件未就绪时返回503 Service Unavailable而非 500 Internal Server ErrorFLINK-25269Avro 生成代码的命名空间改为org.apache.flink.avro.generatedFLINK-25962。Clusters Deploymentjobmanager.sh 的 host/web-ui-port 参数弃用发布说明指出FLINK-28735jobmanager.sh脚本中的host与web-ui-port命令行参数已被弃用应改用对应的动态属性dynamic properties来指定例如通过-Djobmanager.bind-host...、-Drest.address...、-Drest.port...等 option 的方式传入。迁移建议在启动脚本中把原先写死在jobmanager.sh命令行的 host/port统一挪到config.yaml或flink-conf.yaml与-D参数中管理这样同一份集群配置可以在 standalone、YARN 等部署形态间复用。升级时若仍使用旧参数1.16 会给出弃用告警属于软迁移窗口期但后续版本大概率会彻底移除。Table API SQL移除字符串表达式 DSLString Expression DSL发布说明FLINK-26704确认从 Java/Scala/Python 三种 Table API 中移除了此前已弃用的 String 表达式 DSL即形如table.select($id 1, lower($name))的字符串写法。迁移方式改用类型安全的表达式 API。Java 侧使用Expression构建器Expression idPlusOne $id().plus(1); Expression lowerName Functions.lower($name()); table.select(idPlusOne, lowerName);PythonPyFlink侧使用表达式 API 的函数式写法from pyflink.table import expressions as expr table.select(expr.col(id) 1, expr.expr(lower(%s) % name).if_supported()) # 或更推荐 table.select(expr.col(id) 1, expr.lower(expr.col(name)))由于这是直接移除而非弃用所有仍在使用字符串 DSL 的作业会在 1.16 上直接编译失败升级前必须先完成改写。新增Retryable Lookup Join 解决外部维表更新延迟发布说明FLINK-28779介绍了 1.16 的一个亮点能力为同步/异步 Lookup Join 增加可重试查询retryable lookup join用来解决外部维表如缓存、数据库更新延迟导致的查不到维度问题。配置方式通过 HINT 指定重试谓词与重试策略在仓库源码中重试相关 HINT 键定义在 LookupJoinHintOptions.java 中计划侧解析逻辑位于 FlinkHintStrategies.java 与 LookupJoinUtil.java。从源码结构看重试配置支持四个键retry-predicate判定失败类型何种查询失败才需要重试retry-strategy重试策略支持固定间隔FIXED_DELAY与指数退避EXPONENTIAL_DELAY两类fixed-delay固定延迟时间如1smax-attempts最大重试次数。SQL 中的典型用法通过/* LOOKUP(...) */HINT 注入到具体的 Lookup Join 节点上SELECT /* LOOKUP( tabledim, retry-predicateEXCEPTION_AS_FAILURE, retry-strategyFIXED_DELAY, fixed-delay1s, max-attempts3 ) */ d.name, o.* FROM orders AS o JOIN dim FOR SYSTEM_TIME AS OF o.proc_time AS d ON o.id d.id;这些重试参数最终会被序列化进逻辑计划执行节点 CommonExecLookupJoin.java 中定义了retryOptions字段FIELD_NAME_RETRY_OPTIONS流式执行节点 StreamExecLookupJoin.java 在运行时接收该字段并驱动重试。可以推断retryOptions随计划 JSON 一起序列化因此重试行为属于编译期确定的配置升级/恢复 Savepoint 后行为保持一致。新增table.exec.async-lookup.output-mode让异步查询输出模式可配置发布说明FLINK-27622建议当结果不需要严格保序时把新选项table.exec.async-lookup.output-mode设置为ALLOW_UNORDERED在 append-only 流上可显著降低延迟。该选项在仓库中的定义位于 ExecutionConfigOptions.javapublic static final ConfigOptionAsyncOutputMode TABLE_EXEC_ASYNC_LOOKUP_OUTPUT_MODE key(table.exec.async-lookup.output-mode) .enumType(AsyncOutputMode.class) .defaultValue(AsyncOutputMode.ORDERED) .withDescription( Output mode for asynchronous operations which will convert to AsyncDataStream.OutputMode, ORDERED by default. If set to ALLOW_UNORDERED, will attempt to use AsyncDataStream.OutputMode.UNORDERED when it does not affect the correctness of the result, otherwise ORDERED will be still used.);三点关键信息默认值是ORDERED即 1.16 默认行为不变只有显式配置才可能走无序模式枚举取值包含ORDERED、ALLOW_UNORDERED以及显式的UNORDERED底层映射到 DataStream API 的AsyncDataStream.OutputMode描述中明确了安全边界设为ALLOW_UNORDERED时只有在不影响结果正确性的情况下才真正切换为 UNORDERED否则仍回退到 ORDERED——这是该选项比直接用 DataStream 异步模式更保守的设计。与异步 Lookup Join 相关的完整配置族同文件中定义还包括配置项默认值说明table.exec.async-lookup.output-modeORDERED异步查询输出模式1.16 新增table.exec.async-lookup.buffer-capacity100异步 I/O 最大并发数table.exec.async-lookup.timeout3 min异步操作完成超时加固非确定性更新在 changelog 链路中的正确性检测发布说明FLINK-27849提到对于复杂流作业现在可以在运行前检测并提示changelog 管线中存在的非确定性更新non-deterministic updates所导致的潜在正确性问题。这是一个编译期告警性质的加固当优化器发现算子链路上存在可能被重算retraction/retry的非确定性计算如非确定性的 UDF、可能改变结果集顺序的逻辑时会在作业部署前给出警告避免脏数据静默流入下游。升级 1.16 后建议关注作业提交时的告警日志对命中告警的链路做拆分或物化中间结果处理。ConnectorsHive Sink批模式下默认向 Metastore 上报统计信息发布说明FLINK-28883批模式下Hive Sink 现在默认会为写出的表和分区向 Hive Metastore 上报统计信息文件数很多时该过程可能耗时较长可通过table.exec.hive.sink.statistic-auto-gather.enablefalse关闭。在仓库中该选项定义于 HiveOptions.javapublic static final ConfigOptionBoolean TABLE_EXEC_HIVE_SINK_STATISTIC_AUTO_GATHER_ENABLE key(table.exec.hive.sink.statistic-auto-gather.enable) .booleanType() .defaultValue(true) // 1.16 起默认为 true .withDescription( If its true, Flink will gather statistic automatically during writing Hive Table. ... For ORC and Parquet format, numFiles/totalSize/numRows/rawDataSize can be gathered. For other format, only numFiles/totalSize can be gathered. Note: only batch mode supports auto gather statistic, stream mode doesnt support it yet.);源码确认了发布说明未展开的几个细节默认值确认为true且仅批模式生效流模式尚不支持自动收集统计能力与存储格式相关ORC/Parquet 可收集numFiles/totalSize/numRows/rawDataSize四项其他格式只能收集numFiles/totalSize还配套了一个并发选项table.exec.hive.sink.statistic-auto-gather.thread-num默认 3 个线程用于 ORC/Parquet 统计收集写入分区很多时可调大。升级建议批作业若写出的文件数量庞大且下游不依赖 Metastore 统计如靠 Spark 的 file 级统计做优化可在flink-conf.yaml中显式关闭table.exec.hive.sink.statistic-auto-gather.enable: falseHive 版本支持范围收窄发布说明FLINK-27044Flink 不再支持 Hive 1.x、2.1.x 与 2.2.x原因是这些版本已不被 Hive 社区维护。仓库中可以看到目前保留的 SQL connector 产物flink-sql-connector-hive-2.3.10 与 flink-sql-connector-hive-3.1.3即官方预打包的 Hive 版本只剩 2.3.10 与 3.1.3 两档。升级影响若作业使用 Hive 1.x/2.1/2.2 的 MetaStore/HCatalog1.16 上需要自行升级 Hive 环境或从源码构建对应 connector 并自行保证稳定性。Elasticsearch connector 迁移至独立仓库发布说明FLINK-26884Elasticsearch connector 已从 Flink 主仓库复制到独立的 connector 仓库独立版本化维护。1.16 发布周期内两个仓库的产物同时存在但版本号体系不同随 Flink 主线的产物版本为1.16.0外部独立维护的产物版本为3.0.0。官方建议开发者在本发布周期内迁移到后者以对齐 Flink 官方 connector 独立仓库如 Kafka、JDBC 等的版本节奏。注意事项迁移时要同步更换依赖坐标与版本且外部版本3.0.0面向 Flink 1.16 运行环境升级前请核对依赖树中不再混用新旧两套 Elasticsearch connector 类。Pulsar Connectorcursor API 破坏性变更发布说明FLINK-27399列出了 Pulsar connector cursor API 的三处破坏性变更移除CursorPosition#seekPosition()移除StartCursor#seekPosition()StopCursor#shouldStop的返回值由boolean改为StopCondition。升级影响仅影响自行实现或扩展了 Pulsar cursor 接口的用户直接使用 connector 官方实现的作业不受影响自定义 cursor 需要按新接口签名重写。StreamingFileSink 正式标记弃用发布说明FLINK-27188StreamingFileSink被标记为 deprecated取而代之的是自 Flink 1.12 起逐步统一的FileSink。若项目仍在使用StreamingFileSink及其配套的状态化路径策略 API建议升级到FileSink.forRow(...)/FileSink.forBulkFormat(...)的新 API后者在文件滚动、状态恢复语义上已覆盖旧实现。Avro 生成代码命名空间修正发布说明FLINK-25962Flink 生成的 Avro schema 命名空间改为org.apache.flink.avro.generated以兼容 Avro Python SDK此前生成的 schema 命名空间在 Python 侧无法被正确解析。使用 Avro format 且与 Python 生态互通的场景升级后该问题将自动修复Java 侧若对生成类全限定名有硬编码引用需同步调整。AsyncSink 支持可配置 RateLimitingStrategy发布说明FLINK-28487AsyncSinkWriter新增可配置的RateLimitingStrategysink 实现方可以针对特定 sink 定制请求失败时的限流退避行为不指定时默认沿用AIMDRateLimitingStrategy即 Additive-Increase/Multiplicative-Decrease类似 TCP 拥塞控制思路的自适应限流。这条变更主要面向 sink 的开发者基于 AsyncSink API 封装外部系统写入时可以在失败率高的场景下选择更保守或更激进的速率策略而无需 fork 运行时逻辑。Runtime CoordinationMetrics Reporter 按类名加载方式弃用发布说明FLINK-27206通过metrics.reporter.name.class指定 reporter 类名的配置方式被弃用reporter 实现应当提供MetricReporterFactory所有配置迁移到 factory 机制下。特别地若 reporter 从 plugins 目录加载metrics.reporter.name.class将直接不再生效因为插件类加载隔离下无法跨加载器按类名实例化。迁移检查全局检索配置中的metrics.reporter.*.class确认对应 reporter jar 已提供MetricReporterFactory通过META-INF/services/org.apache.flink.metrics.MetricReporterFactorySPI 注册。Datadog reporter 的tags选项弃用发布说明FLINK-29002DatadogReporter的tags选项被弃用改用通用的scope.variables.additional选项由 metrics core 的 scope 机制统一下发额外标签。仓库中 Datadog reporter 位于 flink-metrics-datadog可在其中确认 factory 化的配置入口。REST API未就绪返回 503 而非 500发布说明FLINK-25269当 REST 请求到达但后端组件尚未就绪时1.16 返回503 Service Unavailable此前返回 500。这是一个面向客户端的语义修正监控探针、CI 中等待 Flink 就绪的脚本应以 503 作为还在启动的正常信号做重试而把 500 视为真正异常。Application 模式 HA 时 JobID 改为基于 cluster ID 生成发布说明FLINK-19358开启 HA 的 application mode 下JobID 不再是全零的0000000000...而是基于 cluster ID 生成。升级影响此前依赖application mode 的 JobID 恒为 000...这一隐含约定的外部工具如脚本中写死 JobID、按 JobID 匹配告警/状态路径的逻辑需要改为通过 REST API 动态查询 JobID。Checkpoints引入 Overdraft Buffer 缓解非对齐 Checkpoint 阻塞发布说明FLINK-26762是 1.16 网络栈层面最重要的改进之一为缓解反压期间 subtask 线程被不可中断地阻塞、导致 unaligned checkpoint barrier 无法注入的问题1.16 引入了超支网络缓冲区overdraft buffers概念——subtask 可以在正常配置的 buffer 数量之外额外向 BufferPool 申请默认最多 5 个超支 buffer保证在背压高峰时仍有空间接收 checkpoint barrier 数据。该行为的开关与额度定义在 NettyShuffleEnvironmentOptions.javataskmanager.network.memory.max-overdraft-buffers-per-gate # 默认 5设为 0 可恢复 1.15 旧行为升级建议该变更会略微增加作业内存占用每个 gate 最多 5 个额外 buffer大规模网络 buffer 配置的作业建议在压测环境核对 TaskManager 峰值内存若升级后出现 OOM 或内存压力可显式设置taskmanager.network.memory.max-overdraft-buffers-per-gate: 0恢复 1.15 前的行为代价是 unaligned checkpoint 在强反压下重新暴露阻塞风险该特性与 unaligned checkpoint 配置execution.checkpointing.unaligned.*配合使用效果最佳二者共同构成 1.16 反压场景下 checkpoint 可用性的改进闭环。table.exec.uid.generation修复 1.15.x 非确定性 UID 问题发布说明FLINK-28861揭示了一个影响面较大的兼容性陷阱Flink 1.15.0 与 1.15.1 为从非编译计划non-compiled plans构建的算子生成了非确定性 UID导致同一作业两次部署的算子 UID 不一致从而使基于 UID 的状态恢复/版本升级变得困难甚至不可能。1.16 通过新配置项table.exec.uid.generation修复其默认行为是不再为这类新管线设置 UID如果用户在 1.15.0/1 上已经接受了当时的 UID 行为并需要与之对齐可显式设置table.exec.uid.generation: ALWAYS升级路径建议若当前跑在 1.15.0/1.15.1 且依赖其 UID 做状态恢复升级 1.16 前务必先确认 Savepoint 中算子 UID 与新版本的匹配情况必要时设置ALWAYS保持兼容1.16 之后新建的 SQL 作业使用默认行为即可避免 UID 在不同部署间漂移。PythonPyFlink 1.16 将是最后支持 Python 3.6 的版本发布说明FLINK-28195Python 3.6 的扩展支持已于 2021 年 12 月 23 日结束官方计划PyFlink 1.16 为最后一个支持 Python 3.6 的版本。升级行动项仍使用 Python 3.6 的 PyFlink 环境在升级至 1.17 及以后版本前必须先完成 Python 运行时到 3.7推荐 3.8/3.9的迁移。仓库中 flink-python 模块的依赖声明可以辅助核对各版本兼容矩阵。依赖升级Hadoop 文件系统实现升级至 3.3.2发布说明FLINK-27308Flink 文件系统实现所依赖的 Hadoop 实现升级至 3.3.2带来两方面收益获得 Hadoop 3.3.x 对应的文件系统特性支持客户端侧的 Flink 状态加密HADOOP-13887 对应能力即 KMS 加密的 FS 上读写 Flink 状态时无需依赖服务端透明处理。使用 S3/GCS/Azure/OSS 等通过 hadoop-fs 实现接入的集群升级后可关注对应 filesystem 模块见 flink-filesystems 目录下的各 hadoop 实现子模块带来的行为变化。Kafka Client 升级至 3.1.1发布说明FLINK-28060Kafka connector 默认使用Kafka client 3.1.1。注意事项Kafka client 3.x 要求 broker 端至少为 0.10.0若集群使用较老的 broker或依赖了被旧 client 容忍的特殊行为升级前应在测试环境验证 consumer 组、offset 提交与序列化路径。Hive 2.3 connector 升级至 2.3.9发布说明FLINK-27063Hive 2.3 connector 版本从旧版升级至2.3.9与前面提到的保留 Hive 2.3.10/3.1.3 预打包 connector共同构成 1.16 的 Hive 支持矩阵。PyFlink 依赖版本更新支持 Python 3.9 与 Apple M1发布说明FLINK-25188为支持 Python 3.9 与 M1ARM64架构PyFlink 更新了一组依赖apache-beam2.38.0 arrow5.0.0 pemja0.2.6自建 PyFlink wheel 或维护内部镜像的团队同步升级这三项依赖即可对齐官方行为尤其pemja是 PyFlink 与 JVM 之间对象桥接的核心组件版本需与 Flink 发行版严格匹配。系统资源 Metrics 依赖更新发布说明同时列出系统资源指标相关依赖的升级com.github.oshi:oshi-core:6.1.5 (MIT License) net.java.dev.jna:jna-platform:5.10.0 net.java.dev.jna:jna:5.10.0影响面为cpu.load、内存等系统级 metric 的采集组件常规部署无需操作仅在自打包发行版shaded dist时需要更新依赖树。总结1.16 升级检查清单结合上述各节按必须先做 / 需要评估 / 可选优化三档整理升级检查清单类别事项依据条目必须先做移除 String Expression DSL 用法改用类型安全 Expression APIFLINK-26704必须先做Hive 1.x/2.1.x/2.2.x 用户升级 Hive 环境或改用 2.3.10/3.1.3 connectorFLINK-27044必须先做Python 3.6 用户规划迁移至 3.71.17 起不再支持FLINK-28195必须先做重写自定义 Pulsar cursorseekPosition已移除FLINK-27399需要评估Hive 批 Sink 统计自动收集是否关闭statistic-auto-gather.enableFLINK-28883需要评估1.15.0/1 升级路径的 UID 匹配table.exec.uid.generationFLINK-28861需要评估强反压作业内存余量 vs. overdraft buffer默认 5可设 0 回退FLINK-26762需要评估metrics reporter 从.class配置迁移到MetricReporterFactoryFLINK-27206需要评估REST 客户端将 503 处理为未就绪而非错误FLINK-25269可选优化Append-only 维表查询作业开启table.exec.async-lookup.output-modeALLOW_UNORDEREDFLINK-27622可选优化维表查询不稳定场景启用 Retryable Lookup Join HINTFLINK-28779可选优化jobmanager.sh的 host/web-ui-port 参数迁移为-D动态属性FLINK-28735Flink 1.16 的整体基调是平滑但有牙齿绝大多数变更是新增能力与加固但 String DSL 移除、Pulsar cursor 破坏性变更、Hive 版本收窄和 PyFlink 3.6 的退役四个点是硬断点升级窗口内应优先安排上述检查项的验证。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink 1.16 版本发布说明与升级指南从 1.15 迁移的关键变更全解析Flink 1.16 版本发布说明与升级指南从 1.15 迁移的关键变更全解析 本文围绕 Apache Flink 1.16 的官方 Release Note大数据流处理批处理数据工程Flink 1.12 发布说明深度解读从 1.11 升级的关键变更与兼容性检查清单Flink 1.12 发布说明深度解读从 1.11 升级的关键变更与兼容性检查清单 本指南以 Flink 1.12 官方 Release Notes 为核心大数据流处理批处理数据工程Apache Flink 1.13 发布说明深度解读Failover、SQL 与运行时行为变更全解Apache Flink 1.13 发布说明深度解读Failover、SQL 与运行时行为变更全解 本文基于 Flink 官方发布说明文档 flink 1.1大数据流处理批处理数据工程上一篇AirLLM 深度拆解:4GB 显存下 70B 大模型推理的最小启动配置下一篇终极Emissary-Ingress金丝雀部署教程如何实现零停机应用发布创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考