Kafka2.1+Zeppelin0.8.1+Spark2.4整合报错:找不到Kafka数据源
你大概率踩了一个新手很容易混淆的坑——你用的是Structured Streaming的API,但添加的是Spark Streaming(基于DStream的旧API)的依赖包,这两个是Spark生态里完全独立的组件,依赖包根本不能通用!
核心问题:依赖包选错了
你下载的spark-streaming-kafka-0-10_2.11.jar是给传统Spark Streaming用的,而你要实现的是Structured Streaming(也就是用spark.readStream.format("kafka")这类API),对应的连接器依赖是spark-sql-kafka-0-10_2.11.jar,必须替换成这个才对。
具体解决方案步骤
1. 获取正确的依赖包
Spark 2.4.x对应的Structured Streaming Kafka连接器版本必须和Spark版本严格匹配(比如你用Spark 2.4.8,就选spark-sql-kafka-0-10_2.11:2.4.8),你有两种获取方式:
- 直接从Maven仓库下载对应版本的jar包;
- 通过Zeppelin自动拉取Maven依赖(推荐,能避免版本不匹配问题)。
2. 在Zeppelin中配置正确的依赖
方式一:通过Spark解释器全局配置
- 打开Zeppelin的Interpreter页面,找到Spark解释器;
- 在Dependencies模块点击Add dependency,输入Maven坐标:
org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.x(把x换成你的Spark具体版本,比如2.4.8); - 点击Save,然后重启Spark解释器(不需要重启整个Zeppelin服务)。
方式二:在Notebook中用%dep临时加载
在你的Notebook开头添加这段代码,直接加载正确的依赖:
%dep z.load("org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.8")
注意:不要同时用两种方式加载依赖,可能会导致类加载冲突。
3. 验证依赖是否正确加载
在Zeppelin的Spark段落里执行以下代码,检查是否加载了正确的Kafka连接器jar:
sc.listJars().filter(_.contains("kafka")).foreach(println)
输出里应该能看到spark-sql-kafka-0-10_2.11-2.4.x.jar,而不是spark-streaming-kafka-0-10开头的jar。
4. 测试Structured Streaming读取Kafka
用这段简单的测试代码验证功能是否正常:
val kafkaDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "test") .load() kafkaDF.printSchema()
如果能正常打印出Kafka消息的Schema,说明问题已经解决了。
额外排查点
- 确保本地Kafka服务正常运行,9092端口能正常访问;
- 如果用本地jar包,检查Zeppelin的Spark解释器配置中
spark.jars参数是否包含了正确的jar路径; - 清理掉之前添加的错误依赖(比如Spark jars目录里的
spark-streaming-kafka-0-10jar),避免类加载冲突。
内容的提问来源于stack exchange,提问作者tardis

