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

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包:

    1. 下载对应版本的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等)。
    2. 将所有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包。
  • 验证包配置是否生效
    启动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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 14:55:29