如何处理Apache Beam中WriteToBigQuery返回的失败行并打印?
处理WriteToBigQuery返回的失败行PCollection
我懂你碰到的问题了——beam.io.WriteToBigQuery这个PTransform返回的不是普通的PCollection,而是包含失败行集合的字典,所以直接在后面链式调用| beam.Map(print)肯定会触发AttributeError,毕竟字典没有pipeline属性嘛。
解决方法很简单,你只需要先把WriteToBigQuery的执行结果存到一个变量里,再从这个字典中取出对应失败行的PCollection,之后就可以像处理普通PCollection一样对它做操作了。
直接上代码示例:
# 假设your_input_pcollection是你要写入BigQuery的数据源PCollection write_outputs = your_input_pcollection | beam.io.WriteToBigQuery( 'thijs:thijsset.thijstable', schema=table_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) # 从返回的字典中提取失败行PCollection,然后打印每一行 _ = write_outputs['FailedRows'] | "Print Failed Rows" >> beam.Map(lambda failed_row: print(f"Failed to write row: {failed_row}"))
补充说明:
- 你看到的字典键
'FailedRows'和文档里提到的BigQueryWriteFn.FAILED_ROWS是等价的,后者其实是一个常量,值就是'FailedRows',所以也可以写成write_outputs[beam.io.gcp.bigquery.BigQueryWriteFn.FAILED_ROWS]来获取对应PCollection。 - 很多示例里没展示后续处理,是因为大部分场景下用户只关注成功写入的数据,忽略失败行;但如果需要做失败重试、日志记录或者告警,就可以用这种方式捕获并处理失败行。
内容的提问来源于stack exchange,提问作者Thijs
相关产品推荐
相关产品推荐

