读取含单表头多文件目录时PySpark DataFrame推断Schema报错求助
解决单表头文件目录的PySpark DataFrame创建问题
我完全懂你碰到的这个麻烦——目录里只有一个文件带表头,其余都是纯数据文件,直接读整个目录时,header=True的逻辑会和实际情况冲突导致报错,毕竟Spark默认会把每个文件的第一行都当成表头处理,这显然不是你想要的。下面是一步步的实用解决方案:
1. 从带表头的文件提取正确Schema
先单独读取那个带表头的文件,让Spark帮你推断出准确的Schema,然后把这个Schema保存下来复用:
# 读取带表头的文件,获取已推断好的Schema schema_source_df = spark.read.csv("/sample/flight/flight_delays1.csv", header=True, inferSchema=True) flight_schema = schema_source_df.schema
2. 用预定义Schema读取整个目录
接下来读取整个目录时,必须把header设为False,这样Spark就不会把任何文件的第一行当成表头,而是直接用我们预先拿到的Schema来解析所有数据:
# 使用提取好的Schema读取目录下所有文件 flights = spark.read.csv("/sample/flight/", header=False, schema=flight_schema)
3. 过滤混入的表头行
这时候你会发现,那个带表头的文件的第一行(也就是原表头)被当成了一条数据行混入了DataFrame,所以需要把它过滤掉。你可以用第一个列名来做判断,比如假设Schema的第一个列是FlightDate:
# 过滤掉表头行(判断第一列的值是否等于列名本身) flights = flights.filter(flights["FlightDate"] != "FlightDate")
如果想要更通用的写法(不管列名是什么都能用),可以借助functools.reduce来检查所有列是否都等于对应的列名:
from functools import reduce header_columns = flight_schema.names # 过滤掉所有列值等于列名的行(也就是原表头行) flights = flights.filter(~reduce(lambda x, y: x & y, [flights[col] == col for col in header_columns]))
为什么直接读目录会报错?
当你设置header=True读取目录时,Spark会尝试将每个文件的第一行解析为表头。但目录里其他文件的第一行是真实数据,不是表头,这就会导致Schema推断混乱(比如数据类型不匹配),最终抛出错误。通过先提取正确Schema再统一读取的方式,就能完美避开这个问题。
内容的提问来源于stack exchange,提问作者mmopu
相关产品推荐
相关产品推荐

