Apache Beam写入BigQuery报错:访问写入结果不存在的属性`0`
问题原因
你遇到的AttributeError是因为**beam.io.WriteToBigQuery返回的不是错误PCollection,而是一个WriteResult对象**。你直接将这个对象当作PCollection来应用beam.Map转换,Beam内部会尝试按索引访问它的元素,而WriteResult没有索引属性,因此触发错误。
要获取写入BigQuery时产生的错误记录,需要从WriteResult对象的errors属性中提取对应的PCollection。
解决方法
修改代码,先捕获WriteToBigQuery的返回结果,再通过.errors获取错误PCollection后进行处理:
def poc1(): import apache_beam as beam # Create pipeline. schema = {'fields': [{'name': 'a', 'type': 'STRING', 'mode': 'REQUIRED'}]} pipeline = beam.Pipeline() # 捕获WriteToBigQuery的返回结果(WriteResult对象) write_result = (pipeline | 'Data' >> beam.Create([2, 2]) | 'CreateBrokenData' >> beam.Map(lambda src: {'a': src} if src == 2 else {'a': '2'}) | 'WriteToBigQuery' >> beam.io.WriteToBigQuery( table='dummy_a_table', dataset='PM_INGEST_TEMP', project="bmas-eu-digi-pipe-dev", schema=schema, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, custom_gcs_temp_location='gs://destination_bucket_final' ) ) # 从WriteResult中提取错误PCollection并处理 if write_result.errors: (write_result.errors | 'PrintErrors' >> beam.Map(print)) # 运行管道 pipeline.run().wait_until_finish() if __name__ == '__main__': poc1()
额外注意事项
- 本地运行认证:确保已设置Google Cloud服务账号认证,通过环境变量指定密钥文件路径:
set GOOGLE_APPLICATION_CREDENTIALS="C:\path\to\your-service-account-key.json" - 依赖安装:确认已安装包含GCP组件的Apache Beam:
pip install apache-beam[gcp] - 数据类型匹配:你的schema中字段
a是STRING类型,而输入的是整数2,BigQuery通常会自动转换,但如果需要严格类型控制,可以在CreateBrokenData步骤中将整数转为字符串:beam.Map(lambda src: {'a': str(src)})
内容的提问来源于stack exchange,提问作者ShubhGurukul
相关产品推荐
相关产品推荐

