【发布时间】:2020-03-24 22:00:55
【问题描述】:
在我的管道中,我使用 WriteToBigQuery,如下所示:
| beam.io.WriteToBigQuery(
'thijs:thijsset.thijstable',
schema=table_schema,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED)
这将返回一个字典,如文档中所述,如下所示:
beam.io.WriteToBigQuery PTransform 返回一个字典,其 BigQueryWriteFn.FAILED_ROWS 条目包含所有 写入失败的行。
我如何打印这个 dict 并将它变成一个 pcollection 或者我如何只打印 FAILED_ROWS?
如果我这样做:| "print" >> beam.Map(print)
然后我得到:AttributeError: 'dict' object has no attribute 'pipeline'
我一定已经阅读了一百条管道,但在 WriteToBigQuery 之后我从未见过任何东西。
[编辑] 当我完成管道并将结果存储在变量中时,我有以下内容:
{'FailedRows': <PCollection[WriteToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn).FailedRows] at 0x7f0e0cdcfed0>}
但我不知道如何在这样的管道中使用此结果:
| beam.io.WriteToBigQuery(
'thijs:thijsset.thijstable',
schema=table_schema,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED)
| ['FailedRows'] from previous step
| "print" >> beam.Map(print)
【问题讨论】:
标签: python-3.x google-bigquery google-cloud-dataflow apache-beam