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

Spark3.2.2+Scala2.12读取Kafka流遇Encoder及依赖加载错误求助

问题:Spark 3.2.2 + Scala 2.12 迁移后Kafka流读取报错

我将原本在Spark 2.2 + Scala 2.11.8环境下正常运行的Kafka流读取代码,迁移到Spark 3.2.2 + Scala 2.12.0环境后,构建时出现错误。

原代码

import spark.implicits._
val kafkaStream = spark
   .readStream
   .format("kafka")
   .option("kafka.bootstrap.servers", settings.kafka.brokers)
   .option("startingOffsets", "latest")
   .option("failOnDataLoss", "false")
   .option("subscribe", "serviceproblems")
   .load()

val dataset = kafkaStream.select($"key", $"value").as[(String, String)]
val mapper = new ObjectMapper
mapper.registerModule(new ServiceProblemDeserializerModule())

报错信息

核心错误

could not find implicit value for evidence parameter of type org.apache.spark.sql.Encoder[(String, String)]
[ERROR]       val dataset = kafkaStream.select($"key", $"value").as[(String, String)]

依赖类加载错误

[ERROR] missing or invalid dependency detected while loading class file 'SQLImplicits.class'.
Could not access type Encoder in package org.apache.spark.sql,
because it (or its dependencies) are missing. Check your build definition for
missing or conflicting dependencies. (Re-run with `-Ylog-classpath` to see the problematic classpath.)
A full rebuild may help if 'SQLImplicits.class' was compiled against an incompatible version of org.apache.spark.sql.

[ERROR] missing or invalid dependency detected while loading class file 'LowPrioritySQLImplicits.class'.
Could not access type Encoder in package org.apache.spark.sql,
because it (or its dependencies) are missing. Check your build definition for
missing or conflicting dependencies. (Re-run with `-Ylog-classpath` to see the problematic classpath.)
A full rebuild may help if 'LowPrioritySQLImplicits.class' was compiled against an incompatible version of org.apache.spark.sql.

[ERROR] missing or invalid dependency detected while loading class file 'package.class'.
Could not access type Row in package org.apache.spark.sql,
because it (or its dependencies) are missing. Check your build definition for
missing or conflicting dependencies. (Re-run with `-Ylog-classpath` to see the problematic classpath.)
A full rebuild may help if 'package.class' was compiled against an incompatible version of org.apache.spark.sql.

[ERROR] missing or invalid dependency detected while loading class file 'Dataset.class'.
Could not access type Encoder in package org.apache.spark.sql,
because it (or its dependencies) are missing. Check your build definition for
missing or conflicting dependencies. (Re-run with `-Ylog-classpath` to see the problematic classpath.)
A full rebuild may help if 'Dataset.class' was compiled against an incompatible version of org.apache.spark.sql.

解决方案

  • 对齐依赖版本:确保所有Spark相关依赖(spark-sql、spark-sql-kafka-0-10等)都使用适配Scala 2.12的版本,例如Spark 3.2.2对应依赖的后缀为_2.12,避免混合Scala 2.11和2.12的依赖包。
  • 确认隐式导入有效性:import spark.implicits._必须在SparkSession实例spark创建之后导入,且位置在使用as[(String, String)]之前,保证隐式Encoder能被正确加载。
  • 清理编译产物并重建:执行清理命令(如sbt clean compile或mvn clean install),移除旧环境的编译残留,避免类版本冲突。
  • 匹配Kafka连接器版本:使用与Spark 3.2.2完全兼容的Kafka连接器,即org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.2,版本不匹配会导致类加载异常。

内容的提问来源于stack exchange,提问作者Chandan Gawri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 10:28:18