如何在Dataform中用BigQuery LoadJob处理PubSub消息:加载正确行并获取失败行
解决BigQuery LoadJob错误行中断加载的问题
要实现正确行正常写入、仅捕获失败行且不中断加载,需修改LoadJobConfig的配置,增加错误容忍和错误记录相关参数,同时调整异常处理逻辑:
关键配置修改
在bigquery.LoadJobConfig中添加以下核心参数:
max_bad_records:设置允许跳过的错误行数(比如设为1000,可按需调整),只有当错误行数超过这个阈值时才会触发加载失败,否则会跳过错误行继续加载正确数据ignore_unknown_values:设为True,忽略JSON中表Schema未定义的字段,避免这类场景触发错误bad_records_uri:指定GCS路径(如gs://your-bucket/bad-records/),BigQuery会将错误行写入该路径,方便后续排查分析
修改后的完整代码
try: job_config = bigquery.LoadJobConfig( source_format=bigquery.SourceFormat.NEWLINE_DELIMITED_JSON, schema_update_options=[ bigquery.SchemaUpdateOption.ALLOW_FIELD_ADDITION, bigquery.SchemaUpdateOption.ALLOW_FIELD_RELAXATION ], write_disposition=bigquery.WriteDisposition.WRITE_APPEND, # 新增错误容忍配置 max_bad_records=1000, ignore_unknown_values=True, # 指定错误行存储的GCS路径,需确保BigQuery有该路径的写入权限 bad_records_uri="gs://your-error-log-bucket/pubsub-errors/" ) load_job = self.client.load_table_from_json( json_text, table_id, job_config=job_config, ) load_job.result() # 加载完成后检查是否存在错误行 if load_job.errors: print(f"加载完成,但存在{len(load_job.errors)}条错误行,详情已写入指定GCS路径") # 可在此处将错误信息写入日志表或进行其他自定义处理 except BadRequest as error: # 仅当错误行数超过max_bad_records阈值时才会进入该分支 print(f"加载失败:{error}") # 从load_job.errors中获取具体错误详情 if load_job.errors: for err in load_job.errors: print(f"错误行详情:{err}")
注意事项
- 确保BigQuery服务账号拥有
bad_records_uri对应GCS桶的写入权限 max_bad_records的值需根据业务场景调整,避免因过多错误行导致无效数据大量加载- 若不需要将错误行写入GCS,也可仅通过
load_job.errors获取错误信息,但max_bad_records仍需设置,否则只要存在错误行就会触发BadRequest异常
内容的提问来源于stack exchange,提问作者jereczq22
相关产品推荐
相关产品推荐

