在PySpark(Databricks)中跳过含未知字段异常的JSON坏记录
解决Spark Streaming读取JSON时未知字段报错的方案
针对你遇到的[UNKNOWN_FIELD_EXCEPTION.NEW_FIELDS_IN_RECORD_WITH_FILE_PATH]错误,以下是两种直接可行的解决方法:
方案一:忽略未知字段(跳过多余字段)
在读取阶段添加ignoreUnknownFields选项,让Spark解析JSON时自动忽略不在当前schema中的字段,避免抛出异常:
sent = spark.readStream.format('cloudFiles') \ .option('cloudFiles.format', 'json') \ .option('multiline', 'true') \ .option('cloudFiles.inferColumnTypes', 'true') \ .option('cloudFiles.schemaLocation', checkpoint_path) \ .option('ignoreUnknownFields', 'true') # 新增:忽略未知字段 .load(raw_files) \ .withColumn('load_ts', F.current_timestamp()) \ .writeStream \ .format('delta') \ .option('checkpointLocation', checkpoint_path) \ .trigger(availableNow=True) \ .option('mergeSchema', 'true') \ .toTable(b_write_path)
方案二:自动演进Schema保留新字段
如果你希望保留这些新字段并同步到Delta表中,可以开启Schema自动演进:
sent = spark.readStream.format('cloudFiles') \ .option('cloudFiles.format', 'json') \ .option('multiline', 'true') \ .option('cloudFiles.inferColumnTypes', 'true') \ .option('cloudFiles.schemaLocation', checkpoint_path) \ .option('cloudFiles.schemaEvolutionMode', 'addNewColumns') # 允许添加新字段到Schema .load(raw_files) \ .withColumn('load_ts', F.current_timestamp()) \ .writeStream \ .format('delta') \ .option('checkpointLocation', checkpoint_path) \ .trigger(availableNow=True) \ .option('mergeSchema', 'true') \ .toTable(b_write_path)
注意事项
- 若使用Schema演进,确保
schemaLocation目录下的历史Schema文件允许更新,必要时可清理该目录后重新启动流任务。 mergeSchema=true配置已在写入阶段启用,确保Delta表能同步新增的字段。
内容的提问来源于stack exchange,提问作者TadeG
相关产品推荐
相关产品推荐

