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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 17:47:35