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

Spark循环追加CSV至Delta表遇同列数据类型冲突,求兼容方案

解决跨CSV数据类型不一致导致Delta表追加失败的问题

问题背景

循环读取根目录下各子文件夹的CSV文件,首次加载成功创建Delta表,但后续追加其他子文件夹CSV时失败,报错数据类型不兼容(如columnA首次被识别为StringType,第二次为IntegerType)。单份CSV含约1000个字段、15GB大小,需合并到一张Delta表并保留FolderName溯源字段。

首次执行代码

Schema = 'customer'
TableName= "customerdata"
FolderName = 'subfolder1' # 循环中会动态变更

dfGetCSV = spark.read.format('csv').options(header=True).option('inferSchema', True).load(f"dbfs:/mnt/rootfolder/{FolderName}/*.csv")

dfGetCSV = dfGetCSV.withColumn("FolderName", lit(FolderName)) \
         .withColumn("CREATEDDATE", current_timestamp())

dfGetCSV.write.format("delta").mode("append").option("overwriteSchema", "true").saveAsTable(f"""{Schema}.{TableName}""")

第二次执行代码及报错

Schema = 'customer'
TableName= "customerdata"
FolderName = 'subfolder2' # 读取第二个子文件夹的文件

dfGetCSV = spark.read.format('csv').options(header=True).option('inferSchema', True).load(f"dbfs:/mnt/rootfolder/{FolderName}/*.csv")

dfGetCSV = dfGetCSV.withColumn("FolderName", lit(FolderName)) \
         .withColumn("CREATEDDATE", current_timestamp())

dfGetCSV.write.format("delta").mode("append").option("overwriteSchema", "true").saveAsTable(f"""{Schema}.{TableName}""")

报错信息:Failed to merge fields 'columnA' and 'columnA'. Failed to merge incompatible data types StringType and IntegerType

问题根源:首次加载时部分列全为Null被Spark推断为StringType,后续CSV中这些列有Integer类型值,导致Delta表无法合并不一致的Schema。

解决方案

方案1:预先定义统一Schema(推荐)

针对字段数量多的场景,先整理所有字段的统一数据类型,读取CSV时强制使用该Schema,避免自动推断的差异。

  1. 收集所有CSV的字段列表,根据业务逻辑确定每个字段的最终兼容类型(例如将冲突字段统一设为StringType)。
  2. 定义Spark StructSchema:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType

# 示例Schema,需根据实际1000个字段扩展
unified_schema = StructType([
    StructField("columnA", StringType(), nullable=True),
    StructField("columnB", IntegerType(), nullable=True),
    # 依次添加剩余字段定义
    StructField("FolderName", StringType(), nullable=False),
    StructField("CREATEDDATE", TimestampType(), nullable=False)
])
  1. 读取CSV时指定该Schema:
dfGetCSV = spark.read.format('csv').options(header=True).schema(unified_schema).load(f"dbfs:/mnt/rootfolder/{FolderName}/*.csv")

dfGetCSV = dfGetCSV.withColumn("FolderName", lit(FolderName)) \
         .withColumn("CREATEDDATE", current_timestamp())

dfGetCSV.write.format("delta").mode("append").saveAsTable(f"""{Schema}.{TableName}""")

如果无法提前整理所有字段,可先读取所有CSV的表头合并去重,生成基础Schema后调整冲突字段类型。

方案2:强制所有字段以StringType读取

若业务允许先以字符串存储所有字段,后续再按需转换类型,可直接统一读取为StringType:

# 读取任意一份CSV获取表头
sample_df = spark.read.format('csv').options(header=True).load("dbfs:/mnt/rootfolder/subfolder1/*.csv")
schema_fields = [StructField(col, StringType(), nullable=True) for col in sample_df.columns]
unified_schema = StructType(schema_fields)

# 循环读取时使用该Schema
dfGetCSV = spark.read.format('csv').options(header=True).schema(unified_schema).load(f"dbfs:/mnt/rootfolder/{FolderName}/*.csv")

dfGetCSV = dfGetCSV.withColumn("FolderName", lit(FolderName)) \
         .withColumn("CREATEDDATE", current_timestamp())

dfGetCSV.write.format("delta").mode("append").saveAsTable(f"""{Schema}.{TableName}""")

优势:无需逐个字段定义,适配字段极多的场景;劣势:丢失原始类型信息,后续需额外转换。

方案3:修复现有Delta表Schema后追加

若已创建Delta表,可先修改表的冲突字段类型,再转换数据后追加:

  1. 修改Delta表Schema:
ALTER TABLE customer.customerdata CHANGE COLUMN columnA columnA STRING;
-- 对所有冲突字段执行类似修改
  1. 读取CSV时转换字段类型匹配表Schema:
from pyspark.sql.functions import col

dfGetCSV = spark.read.format('csv').options(header=True).option('inferSchema', True).load(f"dbfs:/mnt/rootfolder/{FolderName}/*.csv")

# 将冲突字段转换为与表Schema一致的类型
dfGetCSV = dfGetCSV.withColumn("columnA", col("columnA").cast(StringType())) \
         .withColumn("FolderName", lit(FolderName)) \
         .withColumn("CREATEDDATE", current_timestamp())

dfGetCSV.write.format("delta").mode("append").saveAsTable(f"""{Schema}.{TableName}""")

适合冲突字段较少的场景,需手动处理所有不一致字段。

内容的提问来源于stack exchange,提问作者user1403789

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 16:45:06