PyFlink消费Kafka Avro数据时AvroRowDeserializationSchema JVM类不存在的解决
解决PyFlink使用AvroRowDeserializationSchema时的ClassNotFoundException问题
你的核心问题是依赖版本不匹配+缺失Avro核心依赖,具体表现:
- Docker中使用Flink 1.16.2,但代码里手动添加了1.17.1版本的flink-avro jar,版本冲突导致类加载失败
- 缺少Avro核心库(
org.apache.avro.io.DatumReader属于Avro核心包),且代码中添加的avro-1.3.3版本过于陈旧,与Flink 1.16.2不兼容 - Docker中已将flink-avro放入lib目录,代码中重复添加不同版本jar,加剧依赖冲突
解决方案步骤
1. 移除代码中手动添加jar的逻辑
删掉代码中以下冲突逻辑,避免版本矛盾:
AVRO_JAR_PATH = f"file://{current_directory}/avro-1.3.3.jar" FLINK_AVRO_JAR_PATH = f"file://{current_directory}/flink-avro-1.17.1.jar" env = StreamExecutionEnvironment.get_execution_environment() env.add_jars(AVRO_JAR_PATH, FLINK_AVRO_JAR_PATH)
2. 修改Dockerfile,补充Avro核心依赖
在Dockerfile的下载connector libraries部分,添加Flink 1.16.2适配的Avro核心jar(适配版本为1.11.0):
# Download connector libraries RUN wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-json/${FLINK_VERSION}/flink-json-${FLINK_VERSION}.jar; \ wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-csv/${FLINK_VERSION}/flink-csv-${FLINK_VERSION}.jar; \ wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-avro/${FLINK_VERSION}/flink-avro-${FLINK_VERSION}.jar; \ wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-sql-avro/${FLINK_VERSION}/flink-sql-avro-${FLINK_VERSION}.jar; \ wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-avro-confluent-registry/${FLINK_VERSION}/flink-avro-confluent-registry-${FLINK_VERSION}.jar; \ wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/flink/flink-connector-jdbc/${FLINK_VERSION}/flink-connector-jdbc-${FLINK_VERSION}.jar; \ # 添加适配Flink 1.16.2的Avro核心依赖 wget -P /opt/flink/lib/ https://repo.maven.apache.org/maven2/org/apache/avro/avro/1.11.0/avro-1.11.0.jar;
3. 重新构建Docker镜像并运行
执行docker build重新构建镜像,确保所有依赖jar都正确下载到Flink的lib目录中,再启动服务测试。
验证逻辑
- Flink会自动加载
/opt/flink/lib目录下的所有jar,无需代码手动添加 - 所有依赖版本统一为Flink 1.16.2适配的版本,彻底避免版本冲突
- 补充的
avro-1.11.0.jar包含所需的org.apache.avro.io.DatumReader类,解决ClassNotFoundException问题
内容的提问来源于stack exchange,提问作者dyusha32741
相关产品推荐
相关产品推荐

