You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 14:30:08