Databricks Autoloader运行成功但未保存数据,求排查帮助
排查Databricks Autoloader无数据输出问题
针对你遇到的Autoloader运行成功但无数据保存的问题,可按以下步骤排查:
1. 检查文件处理状态与路径
- Autoloader会通过
schemaLocation和checkpointLocation记录已处理文件,若之前运行过作业(哪怕无数据),再次运行trigger(once=True)时会跳过已标记的文件。删除这两个目录后重新运行,验证是否能读取新文件:dbutils.fs.rm("/dbfs/FileStore/temp/schema/", recurse=True) dbutils.fs.rm("/dbfs/FileStore/temp/_checkpoint", recurse=True) - 确认输入路径
/dbfs/FileStore/inbound/dsi/data/下的文件:- 文件后缀必须为
.csv,Autoloader仅识别匹配格式的文件; - 文件未被手动移动或修改过,确保是Autoloader未处理过的新文件。
- 文件后缀必须为
2. 验证Schema与数据格式
- 显式指定Schema避免推断错误:Autoloader自动推断Schema可能因数据格式异常(如
age列含非数字值)导致数据过滤,可手动定义Schema强制匹配:from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 匹配你的数据结构[Fname, Lname, age] custom_schema = StructType([ StructField("Fname", StringType(), nullable=True), StructField("Lname", StringType(), nullable=True), StructField("age", IntegerType(), nullable=True) ]) df = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("header", "true") \ .option("delimiter", ",") \ # 显式指定分隔符,避免格式不匹配 .option("cloudFiles.schemaEvolutionMode", "failOnNewColumns") \ .option("cloudFiles.schemaLocation", "/dbfs/FileStore/temp/schema/") \ .schema(custom_schema) \ # 应用自定义Schema .load("/dbfs/FileStore/inbound/dsi/data/") - 检查原始CSV文件:确认无多余空行、格式统一,列数与表头完全匹配。
3. 验证输出目录与数据格式
- Autoloader默认输出为Parquet格式,而非CSV。若你期望CSV输出,需添加格式配置:
df.writeStream.trigger(once=True) \ .option("checkpointLocation","/dbfs/FileStore/temp/_checkpoint") \ .outputMode("append") \ .format("csv") \ # 指定输出格式为CSV .option("header", "true") \ # 保留表头 .start("/dbfs/FileStore/outbound/dsi/output/") \ .awaitTermination() - 检查输出目录权限:确认当前用户对
/dbfs/FileStore/outbound/dsi/output/有读写权限,可通过dbutils.fs.ls("/dbfs/FileStore/outbound/dsi/output/")查看目录内容。
4. 查看流处理日志
在Databricks作业的日志标签页中查看详细运行日志,重点关注:
- 文件读取阶段的警告(如文件无法解析、Schema不匹配);
- 数据写入阶段的权限或路径错误提示。
内容的提问来源于stack exchange,提问作者marie20
相关产品推荐
相关产品推荐

