View a markdown version of this page

升级到 Flink 2.2:完整指南 - Managed Service for Apache Flink

本文属于机器翻译版本。若本译文内容与英语原文存在差异,则一律以英文原文为准。

升级到 Flink 2.2:完整指南

本指南提供了将您的亚马逊阿帕奇托管服务 Flink 应用程序从 Flink 1.x 升级到 Flink 2.2 的分步说明。这是一次重大版本升级,其中包含重大更改,需要仔细规划和测试。

主要版本升级是单向的

升级操作可以在保持状态的情况下将您的应用程序从 Flink 1.x 移动到 2.2,但在 2.2 状态下您无法从 2.2 移回到 1.x。如果您的应用程序在升级后变得不正常,请使用 Rollback API 以最新快照中的原始 1.x 状态返回 1.x 版本。

先决条件

在开始升级之前:

了解您的迁移路径

您的升级体验取决于您的应用程序与 Flink 2.2 的兼容性。了解这些路径有助于你做好适当的准备并设定切合实际的期望。

路径 1:兼容的二进制文件和应用程序状态

预期发生的情况:

  • 调用升级操作

  • 在应用程序状态过渡的情况下完成向 2.2 的迁移:RUNNING→ → UPDATING RUNNING

  • 保留所有应用程序状态,不会丢失数据或进行重新处理

  • 与次要版本迁移相同的体验

最适合:无状态应用程序或使用兼容序列化的应用程序(Avro、兼容的 Protobuf 架构、不带集合的 POJO)

路径 2:二进制不兼容

预期发生的情况:

  • 调用升级操作

  • 操作失败并通过 Operations API 和日志暴露二进制文件不兼容问题

  • 启用自动回滚功能:应用程序无需您的干预即可在几分钟内自动回滚

  • 禁用自动回滚功能:应用程序在不进行数据处理的情况下保持运行状态;您可以手动回滚到旧版本

  • 修复二进制文件后,使用 UpdateApplication API 获得类似于 Path 1 的体验

最适合:使用在 Flink 任务启动期间检测到的已移除 API 的应用程序

路径 3:不兼容的应用程序状态

预期发生的情况:

  • 调用升级操作

  • 移民最初似乎成功了

  • 状态恢复失败后,应用程序会在几秒钟内进入重启循环

  • 通过显示持续重启的 CloudWatch 指标检测故障

  • 手动调用回滚操作

  • 在启动回滚后的几分钟内恢复生产

  • 州移民核您的申请

最适合:状态序列化不兼容的应用程序(带有集合的 POJO,特定状态) Kryo-serialized

注意

强烈建议创建生产应用程序的副本,并在副本上测试以下每个升级阶段,然后再对生产应用程序执行相同的步骤。

第 1 阶段:准备

更新应用程序代码

更新您的应用程序代码以使其与 Flink 2.2 兼容:

  • 在您的或中将 Flink 依赖项更新到版本 2.2.1 pom.xml build.gradle

  • 将连接器依赖项更新为 Flink 2.2 兼容版本(参见)连接器可用性

  • 移除过时的 API 用法

    • 将 DataSet API 替换为 DataStream API 或表 API/SQL

    • 将旧版SourceFunction/SinkFunction替换为 FLIP-27 源代码和 FLIP-143 接收器 API

    • 用 Java API 替换 Scala API 的使用

  • 更新到 Java 17

上传更新的应用程序代码

  • 使用 Flink 2.2 依赖项构建应用程序 JAR

  • 使用与当前 JAR 不同的文件名上传到 Amazon S3(例如,my-app-flink-2.2.jar

  • 记下升级步骤中使用的 S3 存储桶和密钥

第 2 阶段:启用自动回滚

Auto-rollback 允许在 Apache Flink 的亚马逊托管服务在升级失败时自动恢复到以前的版本。

检查自动回滚状态

AWS 管理控制台:

  1. 导航到您的应用程序

  2. 选择配置

  3. 在 “应用程序设置” 下,验证系统回滚是否已启

AWS CLI:

aws kinesisanalyticsv2 describe-application \ --application-name MyApplication \ --query 'ApplicationDetail.ApplicationConfigurationDescription.ApplicationSystemRollbackConfigurationDescription.RollbackEnabled'

启用自动回滚(如果未启用)

aws kinesisanalyticsv2 update-application \ --application-name MyApplication \ --current-application-version-id <version-id> \ --application-configuration-update '{ "ApplicationSystemRollbackConfigurationUpdate": { "RollbackEnabledUpdate": true } }'

第 3 阶段:拍摄快照(可选)

如果为应用程序启用了自动快照,则可以跳过此步骤,否则在升级之前拍摄应用程序快照以保存应用程序的状态。

从正在运行的应用程序拍摄快照

AWS 管理控制台:

  1. 导航到您的应用程序

  2. 选择快照

  3. 选择创建快照

  4. 输入快照名称(例如,pre-flink-2.2-upgrade

  5. 选择 Create(创建)。

AWS CLI:

aws kinesisanalyticsv2 create-application-snapshot \ --application-name MyApplication \ --snapshot-name pre-flink-2.2-upgrade

验证快照创建

aws kinesisanalyticsv2 describe-application-snapshot \ --application-name MyApplication \ --snapshot-name pre-flink-2.2-upgrade

等到开始SnapshotStatusREADY再继续。

第 4 阶段:升级应用程序

您可以使用UpdateApplication操作升级 Flink 应用程序。

您可以通过多种方式调用 UpdateApplication API:

  • 使用 AWS 管理控制台。

    • 转至 AWS 管理控制台上的应用程序页面。

    • 选择配置

    • 选择新的运行时和要从中启动的快照,也称为还原配置。使用最新设置作为还原配置,从最新的快照启动应用程序。指向 Amazon S3 JAR/zip 上新升级的应用程序。

  • 使用动 AWS CLIupdate-application作。

  • 使用 CloudFormation。

    • 更新该RuntimeEnvironment字段。 CloudFormation 以前会删除该应用程序并创建一个新应用程序,这会导致您的快照和其他应用程序历史记录丢失。现在就RuntimeEnvironment地 CloudFormation 更新您的应用程序,不会删除您的应用程序。

  • 使用 S AWS DK。

    • 有关您选择的编程语言,请参阅 SDK 文档。请参阅UpdateApplication

您可以在应用程序处于 RUNNING 状态或应用程序在 READY 状态中停止时执行升级。适用于 Apache Flink 的亚马逊托管服务可验证原始运行时版本和目标运行时版本之间的兼容性。此兼容性检查在您处于状态时运行,StartApplication如果您在RUNNING状态下升级,则会在下次运行此兼容性检查。UpdateApplication READY

从 “运行” 状态升级

aws kinesisanalyticsv2 update-application \ --application-name MyApplication \ --current-application-version-id <version-id> \ --runtime-environment-update FLINK-2_2 \ --application-configuration-update '{ "ApplicationCodeConfigurationUpdate": { "CodeContentUpdate": { "S3ContentLocationUpdate": { "FileKeyUpdate": "my-app-flink-2.2.jar" } } } }'

从 READY 状态升级

aws kinesisanalyticsv2 update-application \ --application-name MyApplication \ --current-application-version-id <version-id> \ --runtime-environment-update FLINK-2_2 \ --application-configuration-update '{ "ApplicationCodeConfigurationUpdate": { "CodeContentUpdate": { "S3ContentLocationUpdate": { "FileKeyUpdate": "my-app-flink-2.2.jar" } } } }'

第 5 阶段:显示器升级

兼容性检查

  • 使用操作 API 检查升级状态。如果存在二进制不兼容或任务启动问题,则升级操作将失败并显示日志。

  • 如果升级操作成功但应用程序陷入重启循环,则表示该状态与新的 Flink 版本不兼容或者更新后的代码存在问题。回顾Flink 2.2 升级的状态兼容性指南如何识别状态不兼容问题。

监控应用程序运行状况

应用程序状态:

  • 应用程序状态应转换:RUNNINGUPDATINGRUNNING

  • 检查应用程序的运行时间。如果是 2.2,则升级操作成功。

  • 如果您的应用程序已启动RUNNING但仍在较旧的运行时上,则会启动自动回滚功能。操作 API 将操作显示为FAILED。检查日志以查找失败的异常。

此外,在以下位置监控这些指标 CloudWatch:

重启指标:

  • numRestarts: 监控是否出现意外重启 — 如果为零且/ uptimenumRestartsrunningTime正在增加,则升级成功。

检查点指标:

  • lastCheckpointDuration: 应类似于升级前的值

  • numberOfFailedCheckpoints: 应保持在 0

第 6 阶段:验证应用程序行为

应用程序在 Flink 2.2 上运行后:

功能验证

  • 验证是否正在从来源读取数据

  • 验证数据是否正在写入接收器

  • 验证业务逻辑产生预期结果

  • 将输出与升级前基准进行比较

性能验证

  • 监控延迟指标(端到端处理时间)

  • 监控吞吐量指标(每秒记录数)

  • 监控检查点持续时间和大小

  • 监控内存和 CPU 利用率

运行 24 小时以上

允许应用程序在生产环境中运行至少 24 小时,以确保:

  • 没有内存泄漏

  • 稳定的检查点行为

  • 没有意外重启

  • 稳定的吞吐量

第 7 阶段:回滚程序

如果升级失败或应用程序正在运行但运行状况不佳,请回滚到以前的版本。

自动回滚

如果启用了自动回滚且启动期间升级失败,则适用于 Apache Flink 的 Amazon 托管服务会自动恢复到以前的版本。

手动回滚

如果应用程序正在运行但运行状况不佳,请使用 RollbackApplication API:

AWS 管理控制台:

  1. 导航到您的应用程序

  2. 选择操作 → 回

  3. 确认回滚

AWS CLI:

aws kinesisanalyticsv2 rollback-application \ --application-name MyApplication \ --current-application-version-id <version-id>

回滚期间会发生什么:

  • 应用程序停止

  • 运行时恢复到之前的 Flink 版本

  • 应用程序代码恢复到以前的 JAR

  • 应用程序从升级最后一次成功拍摄的快照重新启动

重要
  • 你无法在 Flink 1.x 上恢复 Flink 2.2 快照

  • 回滚使用升级前拍摄的快照

  • 升级前务必拍摄快照(第 3 阶段)

后续步骤

有关升级期间的疑问或问题,请参阅Managed Service for Apache Flink 的故障排除或联系 AWS 支持部门。