Databricks组集群读取CSV遇TextFileFormatEdge$.disabled错误,个人集群正常
问题:Databricks组集群(Unity Catalog)读取CSV文件异常
问题现象
- 个人Databricks集群上运行正常的PySpark读取CSV函数,在组集群上返回空DataFrame,且CSV被解析为单一
value列 - 错误栈跟踪:
org.apache.spark.sql.execution.datasources.csv.TextInputCSVDataSource.infer com.databricks.sql.TextFileFormatEdge$.disabled(TextFileFormatEdge.scala:125) org.apache.spark.sql.execution.datasources.csv.CSVDataSource.inferSchema(CSVDataSource.scala:81) ...
- 物理计划显示:
Scan text … Output [1]: [value#...]
怀疑根因
Unity Catalog/组集群出于治理要求,默认禁用了text/CSV的Schema自动推断功能;而个人集群仍允许自动推断,导致未指定显式Schema的read.csv在组集群上静默失败。
待解决问题
- 如何在禁用Schema推断的Databricks组集群(Unity Catalog)上正确读取该CSV文件?
- 应使用哪些选项或定义Schema,才能将文件读取为5列而非单一
value列?
原代码
def add_processing(df_input): # Reference file used to add extra columns based on date ranges df_ref = spark.read.csv( "/project/data/reference_file.csv", # simplified example path sep=";", header=True, ) # Convert string dates to proper DateType df_ref = ( df_ref.withColumn("start_date", F.to_date("start_date", "dd.MM.yyyy")) .withColumn("end_date", F.to_date("end_date", "dd.MM.yyyy")) .toPandas() ) # Build dynamic column names df_ref["label"] = df_ref["label"].str.replace(" ", "_") df_ref["name"] = "cleanpoint_" + df_ref["id"] + "_" + df_ref["label"] # Filter a subset of rows df_subset = df_input.filter(F.col("category").contains("X")) for i in range(len(df_ref["name"])): df_subset = df_subset.withColumn( df_ref["name"][i], F.when( (df_ref["start_date"][i] <= F.to_date(F.col("some_date_column"))) & (df_ref["end_date"][i] > F.to_date(F.col("some_date_column"))), 1 ).otherwise(0), ) df_other = df_input.filter(~F.col("category").contains("X")) df_output = df_subset.unionByName(df_other, allowMissingColumns=True) return df_output.drop("some_date_column")
已尝试无效方法
- 添加
.cache()/.count()触发读取,问题依旧 - 添加
.option("encoding", "UTF-8")无改善 - 替换为
read.format("csv")并指定header、delimiter参数,仍无效
解决方案
核心思路:指定显式Schema
由于组集群禁用了Schema自动推断,必须手动定义CSV对应的Schema,Spark才能正确解析列结构。
步骤1:定义CSV对应的Spark Schema
假设你的参考CSV有5列,结构如下(根据实际列名和类型调整):
from pyspark.sql.types import StructType, StructField, StringType, DateType # 替换为你的CSV实际列名和数据类型 ref_schema = StructType([ StructField("id", StringType(), nullable=False), StructField("label", StringType(), nullable=True), StructField("start_date", StringType(), nullable=True), # 先按字符串读取,后续转DateType StructField("end_date", StringType(), nullable=True), StructField("other_column", StringType(), nullable=True) # 第5列,根据实际情况调整 ])
注意:日期列先按
StringType读取,后续再用to_date转换,避免因格式不匹配导致读取失败。
步骤2:修改读取CSV的代码
在spark.read.csv中添加schema参数,同时确保指定header=True和sep=";":
df_ref = spark.read.csv( "/project/data/reference_file.csv", sep=";", header=True, schema=ref_schema # 添加显式Schema )
步骤3:优化后续逻辑(可选)
可以避免转换为Pandas,全程用PySpark操作,提升性能并减少数据转换开销:
# 替换原Pandas转换部分 df_ref = ( df_ref.withColumn("start_date", F.to_date("start_date", "dd.MM.yyyy")) .withColumn("end_date", F.to_date("end_date", "dd.MM.yyyy")) .withColumn("label", F.regexp_replace("label", " ", "_")) .withColumn("name", F.concat(F.lit("cleanpoint_"), "id", F.lit("_"), "label")) ) # 动态添加列部分改用PySpark的collect()遍历,避免Pandas索引问题 ref_rows = df_ref.select("name", "start_date", "end_date").collect() df_subset = df_input.filter(F.col("category").contains("X")) for row in ref_rows: df_subset = df_subset.withColumn( row.name, F.when( (row.start_date <= F.to_date(F.col("some_date_column"))) & (row.end_date > F.to_date(F.col("some_date_column"))), 1 ).otherwise(0) )
关键注意事项
- Schema必须与CSV的列顺序、列名完全匹配(注意大小写,Spark默认区分大小写)
- 若CSV中有空值,需将对应字段的
nullable设为True - 不要依赖
inferSchema=True,组集群下该参数会被忽略
内容的提问来源于stack exchange,提问作者Codie
相关产品推荐
相关产品推荐

