本文介绍如何将自建集群上基于低版本开源 Flink(1.15 及以下)的 JAR 作业,跨大版本迁移至阿里云实时计算 Flink 全托管。该场景通常涉及 Flink 引擎大版本升级(例如 1.13→1.17 或 1.13→1.20)与运行平台切换(自建集群→全托管服务)的双重改造。
迁移流程概览
依赖改造:将开源 Flink 依赖替换为 VVR(Ververica Runtime)依赖,移除 Scala 后缀,引入 VVR Connector。
代码改造:替换已弃用的 API,将 Scala 导入改为 Java 导入,替换 Connector 中的 shaded 类。
作业部署:将改造后的 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 前仍可正常使用:
SourceFunctionRichSourceFunctionParallelSourceFunction
如果项目严格禁止弃用警告,可在 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);fixedDelayRestart 和 failureRateRestart 同样需通过 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.Map、scala.collection.immutable.List)。如果作业使用了这些类型,需手动注册对应的序列化器。
替换Connector中的Shaded类
实时计算 Flink 全托管提供的 VVR Connector 对 Kafka 客户端类做了 shade 处理(重新打包到不同的命名空间),以避免与用户 classpath 中的 Kafka 客户端版本产生冲突。因此,代码中直接引用的 Kafka 原生类需替换为 shaded 后的类。
Source端
改造前 | 改造后 |
|
|
|
|
Kafka Source 初始化时,使用 .setStartingOffsets(OffsetsInitializer.latest()) 设置起始消费位点。
Sink端
改造前 | 改造后 |
|
|
验证代码改造
代码改造完成后,执行本地编译和单元测试:
mvn clean package -DskipTests=false确认编译通过且无运行时类找不到(ClassNotFoundException)或方法签名不匹配(NoSuchMethodError)等异常。
步骤三:作业部署
作业部署的详细操作请参见JAR作业开发。以下仅说明迁移场景下的关键配置。
上传JAR包
将本地编译通过的业务 JAR 包上传至实时计算控制台。
将 VVR 依赖 JAR 包上传至控制台(仅需上传一次,同一工作空间内的其他作业可直接引用)。
配置并启动作业
在实时计算控制台部署作业,选择 JAR 作业类型。运行模式支持流模式和批模式。
选择与目标 VVR 版本匹配的引擎版本(例如
vvr-8.0.10-jdk11-flink-1.17)。填写以下配置项:
配置项
说明
JAR URI
选择已上传的业务 JAR 包
Entry Point Class
作业入口类的全限定类名
Entry Point Main Arguments
作业运行参数
附加依赖文件
选择已上传的 VVR 依赖 JAR 包
保存配置后启动作业。
验证作业运行
作业启动后,通过以下方式确认迁移成功:
在作业运维页面确认作业状态为运行中,无持续重启。
检查 Checkpoint 是否正常完成。
对比迁移前后的业务数据输出,确认计算结果一致。
异常排查
迁移过程中可能在不同阶段遇到异常。根据异常发生的时机,排查方向和参考文档有所不同。
部署阶段异常
作业部署失败或启动失败,通常由依赖冲突、类找不到或配置错误引起。
排查方式:
查看启动日志。启动日志包含作业启动阶段的详细信息(包括算子运行前的用户日志输出),适用于排查类加载、依赖冲突和配置错误等问题。
检查 JAR 包中是否存在重复的 Flink 依赖(未设置
provided作用域)或缺失的 Connector 依赖。
相关文档:
运行阶段异常
作业启动成功后在运行过程中出现异常或持续重启,通常由数据处理逻辑、资源不足或外部系统连接问题引起。
排查方式:
查看异常日志。通过作业详情页面的异常日志标签页查看异常堆栈信息。
查看运行日志。运行日志分为 Job Manager 日志和 Task Manager 日志:
Job Manager 日志:通过作业详情页面直接查看,包含作业调度和 Checkpoint 相关信息。
Task Manager 日志:按子任务存放。作业运行中时在"运行中 Task Manager"页面查看,作业停止或失败后在"失效 Task Manager"页面查看。
根据运行日志定位 Task Manager 日志:在 Job Manager 日志中找到异常信息,确认发生错误的 Task 编号(例如
taskmanager-1-145中的145),然后在 Task Manager 列表中定位到对应子任务的日志。查看历史任务日志:作业重启后日志从新启动时间开始记录。查看历史运行日志,需在作业日志页面选择对应时间点启动的作业实例。
检查 Checkpoint 是否正常完成。跨大版本迁移后,Checkpoint 配置和状态后端可能需要调整。
相关文档:
数据正确性异常
作业运行正常但输出数据与迁移前不一致,通常由 API 语义变更、时区处理差异或序列化方式变化引起。
排查方式:
对比迁移前后同一时间窗口的输出数据。
检查是否使用了跨版本语义变更的 API(例如窗口触发时机、时间属性处理)。
检查自定义序列化器是否正确注册(参见本文 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秒检查一次新分区如需关闭动态分区检查,将该参数设置为非正数。