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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 21:39:36