PySpark集成Delta Table遇ModuleNotFoundError问题求助
问题:PySpark与Delta Table配合使用时出现ModuleNotFoundError
我正尝试让PySpark与Delta Table配合使用,已执行pip install delta和pip install delta-spark命令。
我的delta.py脚本如下:
from delta.tables import * from pyspark.sql.functions import * deltaTable = DeltaTable.forPath(spark, "/tmp/delta-table") # 给所有偶数ID加100 deltaTable.update( condition = expr("id % 2 == 0"), set = { "id": expr("id + 100") }) # 删除所有偶数ID的数据 deltaTable.delete(condition = expr("id % 2 == 0")) # 合并(Upsert)新数据 newData = spark.range(0, 20) deltaTable.alias("oldData") \ .merge( newData.alias("newData"), "oldData.id = newData.id") \ .whenMatchedUpdate(set = { "id": col("newData.id") }) \ .whenNotMatchedInsert(values = { "id": col("newData.id") }) \ .execute() deltaTable.toDF().show()
使用的spark-submit命令:
spark-submit --packages io.delta:delta-core_2.12:0.7.0 --master local[*] --executor-memory 2g delta.py
运行后出现错误:
:: loading settings :: url = jar:file:/mnt/spark/jars/ivy-2.5.0.jar!/org/apache/ivy/core/settings/ivysettings.xml Ivy Default Cache set to: /home/eugene/.ivy2/cache The jars for the packages stored in: /home/eugene/.ivy2/jars io.delta#delta-core_2.12 added as a dependency :: resolving dependencies :: org.apache.spark#spark-submit-parent-c6184f38-2d95-498c-b711-ead1c4e98cdc;1.0 confs: [default] found io.delta#delta-core_2.12;0.7.0 in central found org.antlr#antlr4;4.7 in central found org.antlr#antlr4-runtime;4.7 in central found org.antlr#antlr-runtime;3.5.2 in central found org.antlr#ST4;4.0.8 in central found org.abego.treelayout#org.abego.treelayout.core;1.0.3 in central found org.glassfish#javax.json;1.0.4 in central found com.ibm.icu#icu4j;58.2 in central :: resolution report :: resolve 1363ms :: artifacts dl 24ms :: modules in use: com.ibm.icu#icu4j;58.2 from central in [default] io.delta#delta-core_2.12;0.7.0 from central in [default] org.abego.treelayout#org.abego.treelayout.core;1.0.3 from central in [default] org.antlr#ST4;4.0.8 from central in [default] org.antlr#antlr-runtime;3.5.2 from central in [default] org.antlr#antlr4;4.7 from central in [default] org.antlr#antlr4-runtime;4.7 from central in [default] org.glassfish#javax.json;1.0.4 from central in [default] --------------------------------------------------------------------- | | modules || artifacts | | conf | number| search|dwnlded|evicted|| number|dwnlded| --------------------------------------------------------------------- | default | 8 | 0 | 0 | 0 || 8 | 0 | --------------------------------------------------------------------- :: retrieving :: org.apache.spark#spark-submit-parent-c6184f38-2d95-498c-b711-ead1c4e98cdc confs: [default] 0 artifacts copied, 8 already retrieved (0kB/16ms) 23/01/28 14:44:02 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable Traceback (most recent call last): File "/home/eugene/dev/pyspark/delta.py", line 1, in <module> from delta.tables import * File "/home/eugene/dev/pyspark/delta.py", line 1, in <module> from delta.tables import * ModuleNotFoundError: No module named 'delta.tables'; 'delta' is not a package 23/01/28 14:44:02 INFO ShutdownHookManager: Shutdown hook called 23/01/28 14:44:02 INFO ShutdownHookManager: Deleting directory /tmp/spark-a8e1196f-83c5-4c81-865e-9522d4d0c056
请问该问题有解决办法吗?
解决方法
1. 重命名脚本文件
你的脚本名为delta.py,会和导入的delta包名冲突——Python会优先加载当前目录下的delta.py,而非安装的delta-spark包。直接把脚本重命名为delta_demo.py或其他不与包名重复的名称即可。
2. 确保Delta Python包正确加载
如果重命名后仍有问题,可在spark-submit命令中指定Delta的Python包路径:
- 先通过
pip show delta-spark查看包的安装位置,找到delta目录并打包:
cd $(pip show delta-spark | grep Location | cut -d ' ' -f 2) zip -r delta.zip delta/
- 然后在
spark-submit中添加--py-files参数:
spark-submit --packages io.delta:delta-core_2.12:0.7.0 --py-files delta.zip --master local[*] --executor-memory 2g delta_demo.py
3. 对齐Delta版本兼容性
你使用的delta-core版本是0.7.0,需确保delta-spark版本与之匹配:
pip install delta-spark==0.7.0
4. 初始化SparkSession时配置Delta扩展
在脚本开头添加SparkSession初始化代码,指定Delta相关配置:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("DeltaDemo") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() from delta.tables import * from pyspark.sql.functions import * # 后续业务代码...
内容的提问来源于stack exchange,提问作者Eugene Goldberg
相关产品推荐
相关产品推荐

