自建低版本Flink JAR作业迁移至实时计算Flink全托管

更新时间:
复制 MD 格式

本文介绍如何将自建集群上基于低版本开源 Flink(1.15 及以下)的 JAR 作业,跨大版本迁移至阿里云实时计算 Flink 全托管。该场景通常涉及 Flink 引擎大版本升级(例如 1.13→1.17 或 1.13→1.20)与运行平台切换(自建集群→全托管服务)的双重改造。

迁移流程概览

  1. 依赖改造:将开源 Flink 依赖替换为 VVR(Ververica Runtime)依赖,移除 Scala 后缀,引入 VVR Connector。

  2. 代码改造:替换已弃用的 API,将 Scala 导入改为 Java 导入,替换 Connector 中的 shaded 类。

  3. 作业部署:将改造后的 JAR 包上传至实时计算控制台,配置并启动作业。

前提条件

  • 已开通阿里云实时计算 Flink 全托管,并创建工作空间。

  • 本地已安装 Maven 3.x。

  • 如果目标版本为 VVR 11.x(对应 Flink 1.20),本地需安装 JDK 11。VVR 8.x 及以下版本使用 JDK 8。

版本选择

VVR(Ververica Runtime)是阿里云实时计算 Flink 全托管使用的引擎,每个 VVR 版本对应一个开源 Flink 版本。推荐的迁移目标版本如下:

VVR 版本

对应 Flink 版本

JDK 要求

说明

VVR 6.0.7

Flink 1.15

JDK 8

改造量最小,适合保守迁移

VVR 8.0.11

Flink 1.17

JDK 8 或 JDK 11

上一大版本的 LTS 版本,推荐优先选择

VVR 11.6

Flink 1.20

JDK 11

最新版本,需完成本文所有代码改造项

版本选择建议:

  • 优先使用上一个大版本的最后一个小版本(即 LTS 版本)。例如当前最新大版本为 VVR 11.x,则优先选择 VVR 8.0.11。

  • 如果业务需要最新版本的特性,使用当前最新大版本的次新小版本。例如 VVR 11.x 的最新小版本为 11.2,则建议使用 11.1。

  • 如果原始作业基于 Flink 1.13 或更低版本,建议迁移至 VVR 8.0.11(Flink 1.17),兼顾稳定性和功能更新。

重要

本文以从 Flink 1.13 迁移至 Flink 1.20(VVR 11.x)为例,涵盖所有改造项。如果迁移到较低的目标版本,部分改造项(如移除 Scala 后缀)可能不适用,请根据目标版本的 Release Notes 按需取舍。

步骤一:依赖改造

替换Connector依赖

将开源 Flink Connector 替换为实时计算提供的 VVR Connector。在 pom.xml 中添加以下依赖,${vvr.version} 的值根据目标版本设置(例如 1.17-vvr-8.0.11):

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mysql</artifactId>
    <version>${vvr.version}</version>
</dependency>
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-kafka</artifactId>
    <version>${vvr.version}</version>
</dependency>
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-hologres</artifactId>
    <version>${vvr.version}</version>
</dependency>

根据实际使用的 Connector 类型按需添加,不必全部引入。

重要

org.apache.flink 组下以 flink- 开头的非 Connector 依赖,作用域必须设置为 provided(即添加 <scope>provided</scope>),否则会引发依赖冲突。

安装二方包

如果项目依赖未发布到 Maven 中央仓库的二方包,需先将其安装到本地 Maven 仓库。以 awdb-java 为例:

mvn install:install-file \
  -Dfile=awdb-java-2.0.0.jar \
  -DgroupId=io.github.aiwen  \
  -DartifactId=awdb-java \
  -Dversion=2.0.0 \
  -Dpackaging=jar

移除Scala后缀

Flink 1.15 起逐步移除对 Scala 的依赖,Flink 1.20 已完全移除。artifact ID 中的 _2.12_2.11 后缀需去除。如果使用 SBT 构建,%% 需改为 %

改造前:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients_2.12</artifactId>
    <version>1.13.5</version>
    <scope>provided</scope>
</dependency>

改造后:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients</artifactId>
    <version>1.20.1</version>
    <scope>provided</scope>
</dependency>

替换Table Planner依赖

自 Flink 1.14 起,Blink Planner 成为默认且唯一的 Planner。独立的 flink-table-planner-blink 模块已弃用,其内容合并至 flink-table-planner

改造前:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-planner-blink_2.12</artifactId>
    <version>1.13.5</version>
    <scope>provided</scope>
</dependency>

改造后:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-planner_2.12</artifactId>
    <version>1.20.1</version>
    <scope>test</scope>
</dependency>

如果项目使用 Table API,还需添加 API Bridge 依赖:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>1.20.1</version>
    <scope>provided</scope>
</dependency>

验证依赖改造

依赖改造完成后,执行以下命令检查是否存在依赖冲突:

mvn dependency:tree -Dverbose | grep conflict

如果输出为空,表示无冲突。如果存在冲突,通过 <exclusions> 排除低版本依赖。

步骤二:代码改造

替换Scala导入为Java导入

Flink 1.20 中,Java 类比对应的 Scala 类功能更完整。Scala API 类将在 Flink 2.0 中移除。

将以下 Scala 导入:

import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment

替换为 Java 导入:

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment

替换弃用API

SourceFunction(可暂缓)

org.apache.flink.streaming.api.functions.source.SourceFunction 已弃用,新的 Source API 提供了更好的并行度控制。但由于迁移改动量较大,以下类在 Flink 2.0 前仍可正常使用:

  • SourceFunction

  • RichSourceFunction

  • ParallelSourceFunction

如果项目严格禁止弃用警告,可在 Maven 编译插件中添加以下参数忽略:

<arg>-Xlint:-deprecation</arg>

RestartStrategies

org.apache.flink.api.common.restartstrategy.RestartStrategies 已弃用。重启策略需通过 Configuration 对象设置,而非直接修改 StreamExecutionEnvironment

改造前:

env.setRestartStrategy(RestartStrategies.noRestart());

改造后:

Configuration config = new Configuration();
config.set(RestartStrategyOptions.RESTART_STRATEGY, "none");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);

fixedDelayRestartfailureRateRestart 同样需通过 Configuration 对象设置。

fromCollection

StreamExecutionEnvironment.fromCollection 已弃用。

改造前:

env.fromCollection(list);

改造后:

env.fromData(list);

toAppendStream

StreamTableEnvironment.toAppendStream 已弃用。

改造前:

tEnv.toAppendStream(table, Row.class);

改造后:

tEnv.toDataStream(table);

Kryo序列化注册方式变更

自定义序列化器的注册入口从 ExecutionConfig 变更为 SerializerConfig

SerializerConfig serializerConfig = env.getConfig().getSerializerConfig();
serializerConfig.addDefaultKryoSerializer(
    SomeClass.class,
    SomeClassCustomSerializer.class
);

Kryo 不原生支持 Scala 单例与工具类(如 scala.collection.immutable.Mapscala.collection.immutable.List)。如果作业使用了这些类型,需手动注册对应的序列化器。

替换Connector中的Shaded

实时计算 Flink 全托管提供的 VVR Connector 对 Kafka 客户端类做了 shade 处理(重新打包到不同的命名空间),以避免与用户 classpath 中的 Kafka 客户端版本产生冲突。因此,代码中直接引用的 Kafka 原生类需替换为 shaded 后的类。

Source

改造前

改造后

org.apache.kafka.clients.consumer.ConsumerConfig

org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.ConsumerConfig

org.apache.kafka.clients.consumer.OffsetResetStrategy

org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.OffsetResetStrategy

Kafka Source 初始化时,使用 .setStartingOffsets(OffsetsInitializer.latest()) 设置起始消费位点。

Sink

改造前

改造后

org.apache.kafka.clients.producer.ProducerConfig

com.ververica.cdc.connectors.shaded.org.apache.kafka.clients.producer.ProducerConfig

验证代码改造

代码改造完成后,执行本地编译和单元测试:

mvn clean package -DskipTests=false

确认编译通过且无运行时类找不到(ClassNotFoundException)或方法签名不匹配(NoSuchMethodError)等异常。

步骤三:作业部署

作业部署的详细操作请参见JAR作业开发。以下仅说明迁移场景下的关键配置。

上传JAR

  1. 将本地编译通过的业务 JAR 包上传至实时计算控制台。

  2. 将 VVR 依赖 JAR 包上传至控制台(仅需上传一次,同一工作空间内的其他作业可直接引用)。

配置并启动作业

  1. 在实时计算控制台部署作业,选择 JAR 作业类型。运行模式支持流模式和批模式。

  2. 选择与目标 VVR 版本匹配的引擎版本(例如 vvr-8.0.10-jdk11-flink-1.17)。

  3. 填写以下配置项:

    配置项

    说明

    JAR URI

    选择已上传的业务 JAR 包

    Entry Point Class

    作业入口类的全限定类名

    Entry Point Main Arguments

    作业运行参数

    附加依赖文件

    选择已上传的 VVR 依赖 JAR 包

  4. 保存配置后启动作业。

验证作业运行

作业启动后,通过以下方式确认迁移成功:

  • 在作业运维页面确认作业状态为运行中,无持续重启。

  • 检查 Checkpoint 是否正常完成。

  • 对比迁移前后的业务数据输出,确认计算结果一致。

异常排查

迁移过程中可能在不同阶段遇到异常。根据异常发生的时机,排查方向和参考文档有所不同。

部署阶段异常

作业部署失败或启动失败,通常由依赖冲突、类找不到或配置错误引起。

排查方式:

  1. 查看启动日志。启动日志包含作业启动阶段的详细信息(包括算子运行前的用户日志输出),适用于排查类加载、依赖冲突和配置错误等问题。

  2. 检查 JAR 包中是否存在重复的 Flink 依赖(未设置 provided 作用域)或缺失的 Connector 依赖。

相关文档:

作业部署与启动常见问题

运行阶段异常

作业启动成功后在运行过程中出现异常或持续重启,通常由数据处理逻辑、资源不足或外部系统连接问题引起。

排查方式:

  1. 查看异常日志。通过作业详情页面的异常日志标签页查看异常堆栈信息。

  2. 查看运行日志。运行日志分为 Job Manager 日志和 Task Manager 日志:

    • Job Manager 日志:通过作业详情页面直接查看,包含作业调度和 Checkpoint 相关信息。

    • Task Manager 日志:按子任务存放。作业运行中时在"运行中 Task Manager"页面查看,作业停止或失败后在"失效 Task Manager"页面查看。

  3. 根据运行日志定位 Task Manager 日志:在 Job Manager 日志中找到异常信息,确认发生错误的 Task 编号(例如 taskmanager-1-145 中的 145),然后在 Task Manager 列表中定位到对应子任务的日志。

  4. 查看历史任务日志:作业重启后日志从新启动时间开始记录。查看历史运行日志,需在作业日志页面选择对应时间点启动的作业实例。

  5. 检查 Checkpoint 是否正常完成。跨大版本迁移后,Checkpoint 配置和状态后端可能需要调整。

相关文档:

数据正确性异常

作业运行正常但输出数据与迁移前不一致,通常由 API 语义变更、时区处理差异或序列化方式变化引起。

排查方式:

  1. 对比迁移前后同一时间窗口的输出数据。

  2. 检查是否使用了跨版本语义变更的 API(例如窗口触发时机、时间属性处理)。

  3. 检查自定义序列化器是否正确注册(参见本文 Kryo 序列化注册章节)。

相关文档:

数据正确性

常见问题

InvalidPidMappingException

异常信息

Caused by: org.apache.flink.kafka.shaded.org.apache.kafka.common.errors.InvalidPidMappingException: The producer attempted to use a producer id which is not currently assigned to its transactional id.

原因分析

当 Kafka Connector 设置 'sink.delivery-guarantee' = 'exactly-once' 时会启用事务写入,存在Transaction ID过多,保存时间较短被回收的情况。

解决方案

该问题的根本原因在于 Kafka 保存 Transaction ID 的时间较短(云 Kafka 默认 15 分钟)。

  • 可通过调整 Kafka 集群参数增加 Transaction ID 缓存时间来解决。该方案需与云消息队列 Kafka 版产品团队确认后方可实施。

  • 不启用事务写入,下游如果做了幂等去重,接受数据重复。

说明

目前不推荐使用事务写入,详情请参见Kafka Connector注意事项

Kafka部分分区位点异常

Flink 作业运行过程中,如果 Kafka Topic 新增了分区,作业需动态识别新分区以避免消息积压。

默认已开启动态分区检查功能,检查间隔为 5 分钟。如需自定义检查间隔,设置 partition.discovery.interval.ms 参数:

KafkaSource.builder()
    .setProperty("partition.discovery.interval.ms", "10000") // 每10秒检查一次新分区

如需关闭动态分区检查,将该参数设置为非正数。