You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 22:21:38