PySpark读取CSV跳过无效首行(Synapse Notebook,不使用Pandas)
PySpark 解决CSV首行无关内容导致的列名识别错误
直接用PySpark可以通过以下几种方式解决,不需要依赖Pandas:
方法一:先读取为文本行过滤后再构建DataFrame
这种方式先以纯文本形式读取所有行,过滤掉无效的首行,再提取表头并构建Schema:
from pyspark.sql import Row from pyspark.sql.types import StructType, StructField, StringType # 以文本格式读取文件 text_rdd = spark.sparkContext.textFile("your_file_path.csv") # 过滤掉首行无效内容和空行 filtered_lines = text_rdd.filter(lambda line: not line.startswith("Report Name : ") and line.strip()) # 提取表头并定义Schema header = filtered_lines.take(1)[0].split(",") schema = StructType([StructField(col.strip(), StringType(), nullable=True) for col in header]) # 处理数据行(跳过表头行) data_rows = filtered_lines.filter(lambda line: line != ",".join(header)) row_rdd = data_rows.map(lambda line: Row(*[x.strip() for x in line.split(",")])) # 创建最终DataFrame df = spark.createDataFrame(row_rdd, schema) df.show()
方法二:先无表头读取再过滤重命名列
先不指定表头读取所有内容,通过行过滤去掉无效行,再提取表头并重命名列:
from pyspark.sql.functions import row_number from pyspark.sql.window import Window # 不带表头读取全部内容 temp_df = spark.read.csv("your_file_path.csv", header=False) # 添加行号用于过滤 window_spec = Window.orderBy("") temp_df = temp_df.withColumn("row_num", row_number().over(window_spec)) # 提取表头(原文件第二行,对应行号2) header = [col for col in temp_df.filter(temp_df.row_num == 2).collect()[0][:-1]] # 过滤掉前两行(无效行+表头行),并删除行号列 data_df = temp_df.filter(temp_df.row_num > 2).drop("row_num") # 重命名列 data_df = data_df.toDF(*header) data_df.show()
注意事项
- 替换
your_file_path.csv为实际文件路径(Synapse中可以用ADLS路径如abfss://container@storageaccount.dfs.core.windows.net/path/file.csv) - 如果CSV中有复杂分隔符或引号包裹内容,可以在
spark.read.csv中添加sep或quote参数适配
内容的提问来源于stack exchange,提问作者data_engineer_eric
相关产品推荐
相关产品推荐

