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

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在组集群上静默失败。

待解决问题

  1. 如何在禁用Schema推断的Databricks组集群(Unity Catalog)上正确读取该CSV文件?
  2. 应使用哪些选项或定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 07:54:52