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

PySpark Databricks中CSV列名被识别为记录数据的解决方法

解决PySpark读取含换行列名CSV的格式异常问题

问题根源是CSV文件的列名存在换行(例如"Total Cases"被拆分为两行),Spark默认仅将第一行识别为表头,后续包含列名剩余部分的行被误判为数据行,导致仅识别出部分列。

方法一:预处理CSV文件修复换行列名

直接修改CSV文件,将换行的列名合并为单行:

  • 若使用本地文本编辑器,打开文件后删除列名中的换行符,例如把Total\nCases改为Total Cases,保存后用原代码读取即可。
  • 若在Databricks环境,可先查看文件内容确认换行位置:
dbutils.fs.head("dbfs:/FileStore/shared_uploads/mahesh2247@gmail.com/Covid_Live.csv")

再通过集群节点的sed命令批量替换换行:

sed -i ':a;N;$!ba;s/\nTotal Cases/Total Cases/g' /dbfs/FileStore/shared_uploads/mahesh2247@gmail.com/Covid_Live.csv

(根据实际列名调整替换规则,比如Total Deaths对应\nTotal Deaths)

方法二:自定义读取逻辑合并表头行

如果无法修改源文件,可通过RDD操作提取并合并表头行,再读取数据:

  1. 提取并合并表头行(假设表头占2行,根据实际情况调整行数):
# 读取文件所有行
text_rdd = spark.sparkContext.textFile("dbfs:/FileStore/shared_uploads/mahesh2247@gmail.com/Covid_Live.csv")
# 获取表头行
header_rows = text_rdd.take(2)
# 合并表头,去除换行符并拆分列名
merged_header = [col.replace('\n', ' ').strip() for col in ','.join(header_rows).split(',') if col.strip()]
  1. 读取数据部分并绑定正确列名:
# 跳过表头行读取数据,转换为DataFrame并设置列名
df_data = spark.read.format("csv")\
    .option("inferschema", "true")\
    .load("dbfs:/FileStore/shared_uploads/mahesh2247@gmail.com/Covid_Live.csv")\
    .rdd.zipWithIndex()\
    .filter(lambda x: x[1] >= 2)\
    .map(lambda x: x[0])\
    .toDF(merged_header)

注意事项

  • 先通过dbutils.fs.head查看文件前几行,确认表头实际占用的行数,避免合并行数错误。
  • 合并表头后需检查列名数量与数据列数是否匹配,防止出现列数不对应问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 20:40:34