Databricks Autoloader读取空Parquet文件报错的解决方案咨询
解决Databricks Autoloader读取空Parquet文件失败的问题
空Parquet文件缺少Schema元数据,导致Autoloader自动推断Schema时触发FAILED_READ_FILE.NO_HINT错误,中断流任务。以下是几种可行的解决思路和代码修改方案:
方案1:复用目标Delta表的Schema(推荐)
如果目标Unity Catalog Delta表已存在,直接使用表的Schema进行读取,避免自动推断。这样即使遇到空文件,也能按照已有Schema处理,不会因无字段信息报错。
修改后的代码片段:
from pyspark.sql.functions import lit, current_timestamp, input_file_name from pyspark.sql.types import StringType import time current_timestamp_col = lit(current_timestamp()) for row in control_df.collect(): if row['FileType'].upper() == 'PARQUET': catalog = catalog schema = row['DestSchema'] table = row['DestTable'] data_loc = f"/Volumes/{landing_zone}/{schema}/{schema}_vol/{table.lower()}" schema_folder = f"{data_container}/{catalog}/{schema}/" checkpoint_loc = f"{schema_folder}/_checkpoint/{table}" table_uc_path = f"{catalog}.{schema}.{table}" print(f"loading from: {data_loc} ... destination: {table_uc_path}") # 处理表不存在的情况,提前初始化 if not spark.catalog.tableExists(table_uc_path): table_and_checkpoint_prep(table_uc_path, checkpoint_loc) # 获取目标表的Schema target_schema = spark.table(table_uc_path).schema df = spark.readStream \ .format("cloudFiles") \ .schema(target_schema) # 替换inferSchema,直接使用目标表Schema .option("cloudFiles.format", "parquet") \ .option("cloudFiles.schemaLocation", checkpoint_loc) \ .option("cloudFiles.schemaEvolutionMode", "rescue") \ .option("cloudFiles.rescuedDataColumn", "_rescued_data") \ .load(data_loc) \ .withColumn("db_IngestDateTime", current_timestamp_col) \ .withColumn("db_LZfilePath", input_file_name()) # 转换datetimeoffset字段 df = cast_datetimeoffset_columns(df) # 写入Delta表 query = df.writeStream \ .format("delta") \ .option("mergeSchema", True) \ .option("checkpointLocation", checkpoint_loc) \ .trigger(availableNow=True) \ .outputMode("append") \ .toTable(table_uc_path) \ query.awaitTermination() print("Autoload Parquet files complete")
方案2:提前过滤空文件
在读取前遍历着陆区文件,过滤掉大小为0的空Parquet文件,只读取非空文件。适合批处理场景(availableNow触发模式)。
修改后的代码片段:
# ... 保留原有变量定义 ... print(f"loading from: {data_loc} ... destination: {table_uc_path}") # 检查并过滤空文件 files = dbutils.fs.ls(data_loc) non_empty_parquet = [f.path for f in files if f.size > 0 and f.path.endswith(".parquet")] if not non_empty_parquet: print(f"No valid non-empty Parquet files found in {data_loc}, skipping this table.") continue # 处理表不存在的情况 if not spark.catalog.tableExists(table_uc_path): table_and_checkpoint_prep(table_uc_path, checkpoint_loc) df = spark.readStream \ .format("cloudFiles") \ .option("inferSchema", "true") \ .option("cloudFiles.format", "parquet") \ .option("cloudFiles.schemaLocation", checkpoint_loc) \ .option("cloudFiles.schemaEvolutionMode", "rescue") \ .option("cloudFiles.rescuedDataColumn", "_rescued_data") \ .load(",".join(non_empty_parquet)) # 加载过滤后的非空文件 .withColumn("db_IngestDateTime", current_timestamp_col) \ .withColumn("db_LZfilePath", input_file_name()) # ... 后续转换和写入逻辑不变 ...
方案3:添加Schema提示兜底
如果无法提前获取目标表Schema,可以通过cloudFiles.schemaHints指定基础字段作为兜底,避免因空文件无Schema报错。例如:
df = spark.readStream \ .format("cloudFiles") \ .option("inferSchema", "true") \ .option("cloudFiles.format", "parquet") \ .option("cloudFiles.schemaLocation", checkpoint_loc) \ .option("cloudFiles.schemaEvolutionMode", "rescue") \ .option("cloudFiles.rescuedDataColumn", "_rescued_data") \ .option("cloudFiles.schemaHints", "id string") # 指定一个兜底字段 .load(data_loc) \ ...
注意:这种方法仅适合允许临时兜底字段的场景,后续需要结合Schema演化逻辑处理真实数据的字段。
内容的提问来源于stack exchange,提问作者Htape
相关产品推荐
相关产品推荐

