PySpark读写CSV/Parquet转Delta报错及相关操作咨询
问题解答与脚本修正
1. 成功创建Delta表的方法
你的报错是Delta Lake与Spark版本不兼容导致的。你使用的Spark 3.5.1需要匹配Delta Lake 3.0.0及以上版本,而当前脚本中用的Delta 2.1.0仅兼容Spark 3.3.x版本。
修正步骤:
- 更新SparkSession配置中的Delta依赖版本为
io.delta:delta-core_2.12:3.0.0(或更高兼容Spark 3.5.x的版本) - 添加Delta Catalog的配置项(Delta 3.x版本要求)
- 确保SparkSession创建后再执行后续操作
修正后的SparkSession代码:
spark = ( SparkSession.builder.master("local[*]") .config("spark.jars.packages", "io.delta:delta-core_2.12:3.0.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .getOrCreate() )
之后执行df_parq.write.format("delta").mode("overwrite").save("output/BankChurners_delta_table")即可成功创建Delta表(添加mode="overwrite"可避免重复创建时的冲突)。
2. CSV/Parquet读取方式推荐
两种方式功能完全等价,只是写法不同:
- 方式I(
format("csv").option().load())更显式,适合新手理解底层逻辑,当需要配置大量自定义参数时更清晰 - 方式II(
spark.read.csv())是Spark提供的便捷封装,代码更简洁
推荐日常使用方式II,除非需要特别多的自定义配置(比如指定分隔符、编码等),此时方式I更易读。
3. DeltaTable.forPath是否正确
是的,DeltaTable.forPath(spark, "/path/to/table")是从已有Delta路径创建DeltaTable对象的标准方法,完全正确。如果Delta表已注册到Spark Catalog,也可以用DeltaTable.forName(spark, "table_name")。
4. 向Delta表追加记录的方法
有两种常用方式:
方式1:直接用DataFrame的append模式写入
适合纯追加新数据(不处理重复或更新):
# 假设new_df是要追加的新数据DataFrame,结构与原Delta表一致 new_df.write.format("delta").mode("append").save("output/BankChurners_delta_table")
方式2:使用DeltaTable的merge操作(Upsert)
如果需要同时处理更新已有记录+插入新记录(即Upsert),可以用Delta的merge API:
deltaTable = DeltaTable.forPath(spark, "output/BankChurners_delta_table") # 假设new_df有唯一标识列(比如CustomerId) deltaTable.alias("old") .merge( new_df.alias("new"), "old.CustomerId = new.CustomerId" # 匹配条件 ) .whenMatchedUpdateAll() # 匹配到则更新所有字段 .whenNotMatchedInsertAll() # 未匹配到则插入新记录 .execute()
5. 脚本中无需导入的模块
以下模块未在脚本中使用,可以删除:
from pyspark import SparkFiles:未用到SparkFiles的功能from pyspark.context import SparkContext:SparkSession会自动创建SparkContext,无需单独导入from pyspark.sql.types import *:脚本中用了inferSchema=True,未手动定义Schema,因此不需要
脚本其他错误修正
- 提前引用spark变量:原脚本中在创建SparkSession前就打印
spark.version,会导致NameError,需将print语句移到SparkSession创建之后 - 缺少sys模块导入:原脚本中打印
sys.executable但未导入sys,需添加import sys - 冗余的PySpark版本打印:
spark.version和spark.sparkContext.version结果一致,无需重复打印
修正后的完整脚本
import sys from pyspark.sql import SparkSession from pyspark.sql.functions import * from delta import * from delta.tables import * # 创建SparkSession spark = ( SparkSession.builder.master("local[*]") .config("spark.jars.packages", "io.delta:delta-core_2.12:3.0.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .getOrCreate() ) # 打印版本信息(移到SparkSession创建后) print("Kernel:", sys.executable) print("Python version:", sys.version) print("Spark version:", spark.version) print("PySpark version:", spark.version) # PySpark版本与Spark版本一致 csv_filename = "dataset/BankChurners.csv" # 读取CSV(用便捷方式) df_csv = spark.read.csv(csv_filename, header=True, inferSchema=True) # 保存为Parquet df_csv.write.parquet("output/BankChurners_spark.parquet", mode="overwrite") # 读取Parquet df_parq = spark.read.parquet("output/BankChurners_spark.parquet") # 保存为Delta表 df_parq.write.format("delta").mode("overwrite").save("output/BankChurners_delta_table") # 获取DeltaTable对象 deltaTable = DeltaTable.forPath(spark, "output/BankChurners_delta_table") # 示例:追加一条记录 # 先构造一条与原表结构一致的新数据 new_data = spark.createDataFrame( [("12345", "John Doe", 30, 1000.0)], # 替换为实际表的字段值 schema=df_parq.schema ) # 追加写入 new_data.write.format("delta").mode("append").save("output/BankChurners_delta_table") # 验证追加结果 deltaTable.toDF().show()
内容的提问来源于stack exchange,提问作者coredump
相关产品推荐
相关产品推荐

