Spark Delta Table依赖无法解析问题求助
解决无网络Spark集群中Delta Table依赖解析失败问题
问题分析
- 网络限制导致依赖拉取失败:提交命令中使用
--packages参数会触发Spark从Maven中央仓库拉取依赖,但服务器无网络访问权限,因此出现连接超时错误。 - 重复依赖配置冲突:代码中同时配置了
spark.jars和spark.jars.packages,提交命令又重复使用--packages和--jars,导致依赖解析逻辑混乱。 - 依赖完整性缺失:Delta Lake需要多个配套jar(如
delta-storage),仅提供delta-sparkjar会导致依赖不完整。
解决方案
步骤1:准备完整的Delta本地依赖包
在有网络的机器上下载Delta Lake 3.2.0的所有依赖jar,至少包含:
delta-spark_2.12-3.2.0.jardelta-storage-3.2.0.jar
可通过Maven命令一键下载所有依赖:
mvn dependency:copy-dependencies -Dartifact=io.delta:delta-spark_2.12:3.2.0 -DoutputDirectory=./delta-jars
将下载好的所有jar上传到集群所有节点的同一目录,比如/opt/jnpm/delta-jars/。
步骤2:修改代码移除冗余依赖配置
删除代码中spark.jars.packages的配置,避免和提交命令冲突:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, DecimalType, TimestampType, LongType from delta import configure_spark_with_delta_pip import os import sys print(sys.executable) os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable os.environ['PYSPARK_PYTHON'] = sys.executable builder = SparkSession.builder \ .appName("DELTA TABLE EXAMPLE") \ .enableHiveSupport() \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ builder = configure_spark_with_delta_pip(builder) spark = builder.getOrCreate() # Schema定义 schema: StructType = StructType([ StructField("ne_id", LongType(), True), StructField("init_time", StringType(), True), StructField("sum", DecimalType(12,2), True), StructField("total", DecimalType(12,2), True), StructField("avg", DecimalType(12,2), True), StructField("partition_date", StringType(), True)]) df = spark.read.format("csv").option("header", "true").schema(schema).load("/user/spark/test1.csv") df.show() print(df.printSchema()) df.write.format("delta").partitionBy("partition_date").mode("overwrite").save("/user/spark/examples/delta/test")
步骤3:修改提交命令,仅使用本地jar
移除--packages参数,用--jars指定所有本地Delta依赖jar(多个jar用逗号分隔):
spark-submit --master spark://sparkhost:7077 \ --deploy-mode client \ --executor-memory 1g \ --executor-cores 1 \ --num-executors 1 \ --total-executor-cores 1 \ --jars /opt/jnpm/delta-jars/delta-spark_2.12-3.2.0.jar,/opt/jnpm/delta-jars/delta-storage-3.2.0.jar \ --conf "spark.pyspark.python=/opt/jnpm/spark/venv/bin/python" \ Test_Hdfs.py
额外注意事项
- 确保所有集群节点的jar路径一致,避免Executor节点找不到依赖。
- 确认Spark版本与Delta版本兼容:Spark 3.5.0支持Delta Lake 2.4.0至3.2.0版本,当前使用的3.2.0版本是兼容的。
- 虚拟环境中需提前安装对应版本的Delta Python包:
pip install delta-spark==3.2.0,避免Python端依赖缺失。
内容的提问来源于stack exchange,提问作者vivekdesai
相关产品推荐
相关产品推荐

