Databricks中Spark读取列数不同CSV合并为DataFrame列缺失问题
问题场景
- 已在Azure平台配置Blob存储容器,需要将容器内所有.csv格式文件加载到同一个Spark DataFrame中
- 所有文件的前两列固定为
name和time,要求实现的处理逻辑:- 将
time列转换为datetime类型 - 根据文件名生成新的
id列,并将该列调整为DataFrame第一列
- 将
- 除固定列外,其余扩展列遵循统一命名规则(LVxx/CHxx格式),但不同文件包含的扩展列数量存在差异:部分文件仅包含LV01~LV03共3个扩展列,部分文件扩展列到LV15,还有部分文件扩展列最高到LV25及以上。
单文件结构示例:
| id | name | time | LV01 | LV02 | LV03 |
|---|---|---|---|---|---|
| abc | name1 | 01/01/1900 01:00:00 | 47.96 | 23.10 | 43.00 |
现有实现代码
from pyspark.sql.functions import * file_location = 'dbfs:/mnt/<container>/<foldername>' file_type = "csv" infer_schema = "true" first_row_is_header = "true" delimiter = "," df3 = spark.read.format(file_type) \ .option("inferSchema", infer_schema) \ .option("header", first_row_is_header) \ .option("sep", delimiter) \ .load(file_location) \ .withColumn("id",substring(input_file_name(), 45, 3)) # 从文件名提取id字段 df3 = df3.select([df3.columns[-1]] + df3.columns[:-1]) # 将id列移动到第一列 df3 = df3.withColumn("time", (col('time')/1000000000)) # 纳秒转秒,Spark原生不支持纳秒级时间戳转换 df3 = df3.withColumn("time",from_unixtime(col('time'))) # 秒级时间戳转datetime格式 df3.show()
遇到的问题
- 经校验
id列生成结果符合预期,文件已被成功加载,但DataFrame列数存在异常:仅能获取到最高到CH14的列,已知部分文件包含最高到CH25的列,存在明显的列读取缺失问题。 - 后续需要实现类似pandas
melt的长表转换逻辑,对应pandas实现代码如下:
cols = df.iloc[:,3:-1:] col_names = list(cols.columns.values) col_names df_long = df.melt(id_vars=['id','name', 'time'], var_name='channel',value_vars=col_names, value_name='value') df_long.head()
- 待明确的核心问题:
- Spark加载多文件时是否仅会加载所有文件共有的一致列?
- 如何调整代码才能正确读取全部列完成合并,支撑后续的长表转换操作?
解决方案
核心问题原因
Spark加载多CSV文件时不存在默认只读公共列的逻辑,列缺失的根源是:开启inferSchema=true时,Spark默认只会采样路径下的部分文件推导表结构和字段类型,未被采样到的、包含更多扩展列的文件,多出的字段会被直接丢弃,这和你观察到的“只能读到CH14、读不到更高序号列”的现象完全吻合。
修复列读取不全问题
两种可落地方案,优先选择第一种稳定性更高:
- 方案1:全量收集表头构建完整Schema,显式指定Schema读取(推荐,不会受采样规则影响)
核心逻辑是先遍历所有CSV文件提取表头,合并得到所有列的并集,再基于全量Schema读取数据,从根源避免采样漏列。示例代码:from pyspark.sql.types import StructType, StructField, StringType, DoubleType # 第一步:遍历目标路径下所有csv文件,收集全量列名 all_columns = set() files = dbutils.fs.ls(file_location) csv_files = [f.path for f in files if f.path.endswith('.csv')] for file in csv_files: # 读取单个文件第一行提取表头 header = spark.read.text(file).limit(1).collect()[0][0] cols = [c.strip() for c in header.split(delimiter)] all_columns.update(cols) # 第二步:构建完整Schema,固定列name、time先设为StringType后续再转换,扩展列统一设为DoubleType fixed_cols = [ StructField("name", StringType(), nullable=True), StructField("time", StringType(), nullable=True) ] extend_cols = [ StructField(col, DoubleType(), nullable=True) for col in sorted(list(all_columns - {"name", "time"})) ] full_schema = StructType(fixed_cols + extend_cols) # 第三步:基于完整Schema读取所有文件,关闭自动推Schema df3 = spark.read.format(file_type) \ .option("header", first_row_is_header) \ .option("sep", delimiter) \ .schema(full_schema) \ .load(file_location) \ .withColumn("id",substring(input_file_name(), 45, 3)) # 调整id列到第一列 df3 = df3.select(["id"] + [c for c in df3.columns if c != "id"]) - 方案2:调整采样配置,强制Spark扫描全量文件推导Schema
适合文件总量不大的场景,不需要手动遍历文件,修改读取参数即可:df3 = spark.read.format(file_type) \ .option("inferSchema", "true") \ .option("header", first_row_is_header) \ .option("sep", delimiter) \ .option("samplingRatio", 1.0) # 采样比例设为1即扫描全量文件推Schema .load(file_location) \ .withColumn("id",substring(input_file_name(), 45, 3))
后续长表转换实现
确认所有列读取完整后,直接使用Spark内置的melt函数即可实现和pandas melt完全一致的效果,不需要硬编码列名:
# 自动提取所有扩展列(排除固定列id、name、time) extend_cols = [c for c in df3.columns if c not in ["id", "name", "time"]] # 转换为长表 df_long = df3.melt( ids=["id", "name", "time"], values=extend_cols, variableColumnName="channel", valueColumnName="value" ) # 完成time列的类型转换 df_long = df_long.withColumn("time", (col("time").cast("long")/1000000000)) df_long = df_long.withColumn("time", from_unixtime(col("time")))
内容的提问来源于stack exchange,提问作者JGW
相关产品推荐
相关产品推荐

