如何使用Databricks Autoloader跳过CSV文件开头的空行?
跳过CSV开头空行的Databricks Autoloader解决方案
针对CSV文件开头存在大量空行的场景,有两种简便处理方法:
方法一:固定跳过指定行数(已知空行数量时)
如果明确知道开头空行的具体数量,直接在Autoloader的CSV配置中添加skipRows参数,指定要跳过的行数。示例代码:
df = spark.readStream \ .format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("skipRows", "10") # 替换为实际空行数量 .option("header", "true") # 若CSV含表头,需开启此选项 .option("cloudFiles.schemaLocation", "/dbfs/path/to/schema") \ .load("/dbfs/path/to/source/csv") df.writeStream \ .format("delta") \ .option("checkpointLocation", "/dbfs/path/to/checkpoint") \ .table("your_target_delta_table")
方法二:动态过滤全空行(空行数量不固定时)
如果不同文件的空行数量不一致,或无法确定具体行数,可以读取所有行后,过滤掉所有字段均为空的行。示例代码:
from pyspark.sql.functions import col, trim, length, concat_ws df = spark.readStream \ .format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("header", "true") \ .option("cloudFiles.schemaLocation", "/dbfs/path/to/schema") \ .load("/dbfs/path/to/source/csv") # 过滤所有字段拼接后为空的行 filtered_df = df.filter(length(trim(concat_ws(" ", *df.columns))) > 0) filtered_df.writeStream \ .format("delta") \ .option("checkpointLocation", "/dbfs/path/to/checkpoint") \ .table("your_target_delta_table")
注意事项
- 若CSV包含表头,需确保
header选项设为true,方法一中的skipRows参数仅针对表头之前的空行。 - 动态过滤方法会遍历每行数据,对超大文件的性能有轻微影响,但适配性更强。
内容的提问来源于stack exchange,提问作者bradmuzza
相关产品推荐
相关产品推荐

