在PySpark中使用Delta Lake时,如何配置Kafka作为依赖?
解决Spark+Delta Lake集成Kafka时找不到数据源的问题
核心问题是Kafka依赖包版本与PySpark版本不匹配:你使用的PySpark是3.3.0,但指定的Kafka包版本是3.3.1,Spark的Kafka集成包版本必须和PySpark主版本严格一致。
修正后的配置代码
import pyspark from delta import * # 匹配PySpark 3.3.0的Kafka依赖版本 packages = [ "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0", ] builder = pyspark.sql.SparkSession.builder.appName("MyApp") \ .config("spark.jars.packages", ",".join(packages)) \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") spark = configure_spark_with_delta_pip(builder).getOrCreate()
额外排查点
- 确认Scala版本后缀:
_2.12要和PySpark内置的Scala版本一致,PySpark 3.3.x默认使用Scala 2.12,这个配置是正确的。 - 若依赖拉取失败,可手动下载对应版本的Kafka jar包,放到Spark安装目录的
jars文件夹下,跳过在线拉取环节。 - 验证配置是否生效:创建SparkSession后执行
print(spark.conf.get("spark.jars.packages")),确认Kafka依赖包的配置已被正确加载。
内容的提问来源于stack exchange,提问作者Gustavo Puma
相关产品推荐
相关产品推荐

