Spark Streaming连接Kafka报AnalysisException错误的解决咨询
解决Spark Streaming读取Kafka时的AnalysisException(找不到kafka数据源)
问题场景
虚拟机(VMWare)运行Spark 3.5.0 + Kafka 3.5,笔记本与虚拟机处于同一网络:
- 笔记本Python生产者可正常向虚拟机Kafka发送消息
- 相同脚本在虚拟机Jupyter中运行正常,但笔记本本地运行报错:
AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of Structured Streaming + Kafka Integration Guide.
代码示例
from pyspark.sql import SparkSession from pyspark import SparkContext, SparkConf from pyspark.sql.functions import * from pyspark.sql.types import * import os os.environ["JAVA_HOME"] = 'C:\Program Files\Java\jdk-1.8' KAFKA_BOOTSTRAP_SERVERS = "192.168.xxx.xxx:9092" KAFKA_TOPIC = "test" spark = SparkSession.builder.appName("test_stream").master("local[*]").config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0").getOrCreate() df = spark.readStream.format("kafka") \ .option("bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) \ .option("subscribe", KAFKA_TOPIC) \ .load() df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
解决方法
版本兼容性检查
确认本地Spark的Scala版本与Kafka连接器版本匹配:- 连接器包名
spark-sql-kafka-0-10_2.12:3.5.0中的_2.12对应Scala 2.12版本,本地Spark必须是基于Scala 2.12编译的。 - 验证方式:启动
spark-shell,执行scala.util.Properties.versionString查看Scala版本;或检查Spark安装包文件名(如spark-3.5.0-bin-scala2.12.tgz对应Scala 2.12)。 - 如果本地Spark是Scala 2.13版本,需将连接器包改为
org.apache.spark:spark-sql-kafka-0-10_2.13:3.5.0。
- 连接器包名
手动加载Kafka连接器jar包
若自动下载包(spark.jars.packages)因网络问题失败,手动下载并放置jar包:- 下载对应版本的
spark-sql-kafka-0-10_2.12-3.5.0.jar及依赖包(如kafka-clients-3.5.0.jar、spark-token-provider-kafka-0-10_2.12-3.5.0.jar等)。 - 将所有jar包放到本地Spark安装目录的
jars文件夹下;或在SparkSession配置中指定本地jar路径:spark = SparkSession.builder.appName("test_stream")\ .master("local[*]")\ .config("spark.jars", "file:///path/to/spark-sql-kafka-0-10_2.12-3.5.0.jar,file:///path/to/kafka-clients-3.5.0.jar")\ .getOrCreate()
- 下载对应版本的
修复环境变量配置
- Windows下
JAVA_HOME路径需用原始字符串或双反斜杠,避免转义错误:os.environ["JAVA_HOME"] = r'C:\Program Files\Java\jdk-1.8' # 或 os.environ["JAVA_HOME"] = 'C:\\Program Files\\Java\\jdk-1.8' - 确认本地已设置
SPARK_HOME环境变量,指向Spark安装目录,确保SparkSession能正确加载核心jar包。
- Windows下
验证包配置是否生效
启动SparkSession后,执行以下代码检查包配置:print(spark.conf.get("spark.jars.packages"))若未正确返回配置的包名,添加maven镜像配置:
spark = SparkSession.builder.appName("test_stream")\ .master("local[*]")\ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0")\ .config("spark.jars.repositories", "https://repo1.maven.org/maven2")\ .getOrCreate()
内容的提问来源于stack exchange,提问作者Amatuer
相关产品推荐
相关产品推荐

