Spark3.1.2环境下导入pyspark.streaming.kafka模块出现ModuleNotFoundError求助
解决PySpark中
pyspark.streaming.kafka模块找不到的问题 问题背景
环境配置:
SPARK_VERSION = '3.1.2' SCALA_VERSION = '2.12'
执行以下代码时出现模块找不到错误:
import findspark findspark.add_packages(['org.apache.spark:spark-sql-kafka-0-10_' + SCALA_VERSION + ':' + SPARK_VERSION ]) findspark.init() from pyspark import SparkContext, SparkConf import sys import time from pyspark.context import SparkContext from pyspark import SparkContext, SparkConf from pyspark.streaming import StreamingContext from pyspark.streaming.kafka import KafkaUtils
错误信息:
ModuleNotFoundError Traceback (most recent call last) /tmp/ipykernel_567977/2450515063.py in <module> 4 from pyspark import SparkContext, SparkConf 5 from pyspark.streaming import StreamingContext ----> 6 from pyspark.streaming.kafka import KafkaUtils ModuleNotFoundError: No module named 'pyspark.streaming.kafka'
解决方法
1. 修正依赖包加载
你当前添加的是spark-sql-kafka-0-10包,但pyspark.streaming.kafka属于Spark Streaming的旧Kafka API,对应的依赖包是spark-streaming-kafka-0-10。修改findspark.add_packages参数:
findspark.add_packages([ 'org.apache.spark:spark-streaming-kafka-0-10_' + SCALA_VERSION + ':' + SPARK_VERSION, 'org.apache.spark:spark-sql-kafka-0-10_' + SCALA_VERSION + ':' + SPARK_VERSION # 若需SQL Kafka支持可保留 ])
2. 改用Structured Streaming(推荐)
Spark 2.3+官方推荐使用Structured Streaming处理Kafka流,API更简洁且功能更完善,无需依赖旧的pyspark.streaming.kafka模块。示例代码:
import findspark findspark.add_packages(['org.apache.spark:spark-sql-kafka-0-10_' + SCALA_VERSION + ':' + SPARK_VERSION ]) findspark.init() from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder \ .appName("KafkaStructuredStreaming") \ .getOrCreate() # 读取Kafka流 df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "your-kafka-broker:9092") \ .option("subscribe", "your-topic") \ .load() # 转换数据格式 df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") # 启动流查询并等待结束 query = df.writeStream \ .outputMode("append") \ .format("console") \ .start() query.awaitTermination()
3. 确保依赖版本匹配
Spark 3.x中旧的pyspark.streaming.kafka模块仍存在,但需保证依赖包版本与Spark版本完全一致(你的Spark 3.1.2对应Scala 2.12是正确的)。若本地Spark未预装该依赖,提交作业时可通过--packages指定:
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-10_2.12:3.1.2 your_script.py
内容的提问来源于stack exchange,提问作者Chanchal Jain
相关产品推荐
相关产品推荐

