【问题标题】:Monitoring WriteToBigQuery监控 WriteToBigQuery
【发布时间】: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


    【解决方案1】:

    处理无效输入的死信是 Beam/Dataflow 的常见用法,可与 Java 和 Python SDK 一起使用,但后者的示例并不多。

    假设我们有一些虚拟输入数据,其中包含 10 行好行和一个不符合表架构的坏行:

    schema = "index:INTEGER,event:STRING"
    
    data = ['{0},good_line_{1}'.format(i + 1, i + 1) for i in range(10)]
    data.append('this is a bad row')
    

    然后,我要做的是为写入结果命名(在这种情况下为events):

    events = (p
        | "Create data" >> beam.Create(data)
        | "CSV to dict" >> beam.ParDo(CsvToDictFn())
        | "Write to BigQuery" >> beam.io.gcp.bigquery.WriteToBigQuery(
            "{0}:dataflow_test.good_lines".format(PROJECT),
            schema=schema,
        )
     )
    

    然后访问FAILED_ROWS侧输出:

    (events[beam.io.gcp.bigquery.BigQueryWriteFn.FAILED_ROWS]
        | "Bad lines" >> beam.io.textio.WriteToText("error_log.txt"))
    

    这适用于 DirectRunner,并将好的行写入 BigQuery:

    还有一个坏的到本地文件:

    $ cat error_log.txt-00000-of-00001 
    ('PROJECT_ID:dataflow_test.good_lines', {'index': 'this is a bad row'})
    

    如果您使用DataflowRunner 运行它,您将需要一些额外的标志。如果您遇到 TypeError: 'PDone' object has no attribute '__getitem__' 错误,则需要添加 --experiments=use_beam_bq_sink 才能使用新的 BigQuery 接收器。

    如果您收到 KeyError: 'FailedRows',那是因为新接收器将 default 为批处理管道加载 BigQuery 作业:

    STREAMING_INSERTS、FILE_LOADS 或 DEFAULT。加载介绍 数据到 BigQuery:https://cloud.google.com/bigquery/docs/loading-data。 DEFAULT 将在 Streaming 管道上使用 STREAMING_INSERTS 和 批处理管道上的 FILE_LOADS。

    您可以通过在 WriteToBigQuery 中指定 method='STREAMING_INSERTS' 来覆盖该行为:

    DirectRunnerDataflowRunner here 的完整代码。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-22
      • 1970-01-01
      • 1970-01-01
      • 2013-02-19
      相关资源
      最近更新 更多