【问题标题】:WriteToBigQuery with dynamic destinations具有动态目标的 WriteToBigQuery
【发布时间】:2020-11-20 12:08:12
【问题描述】:

我正在开发一个 Apache Beam 管道,该管道从 pub/sub 读取一堆事件,然后根据事件类型将它们写入单独的 BigQuery 表中。

我知道WriteToBigQuery 支持动态目标,但我的问题是目标是从从事件中读取的数据派生的。例如: 一个事件看起来像

{
 "object_id": 123,
 ... some metadata,
 "object_data": {object related info}
}

应写入 BigQuery 表的数据位于事件的 object_data 键下,但表名称源自元数据中的其他字段。 我尝试使用侧输入参数,但问题是因为每个事件都可以有不同的目的地,所以侧输入不会相应地更新。代码如下:

class DumpToBigQuery(PTransform):

    def _choose_table(self, element, table_names):
        # table_names = {"table_name": "project_name.dataset.table_name}
        table_name = table_names["table_name"]
        return table_name

    def expand(self, pcoll):
        events = (
            pcoll
            | "GroupByObjectType" >> Map(lambda e: (e["object_type"], e))
            | "Window"
            >> WindowInto(
                windowfn=FixedWindows(self.window_interval_seconds)
            )
            | "GroupByKey" >> GroupByKey()
            | "KeepLastEventOnly" >> ParDo(WillTakeLatestEventForKey()
        )

        table_name = events | Map(lambda e: ["table_name", f"{self.project}:{self.dataset}.{e[0]}"])
        table_names_dct = AsDict(table_name)

        events_to_write = events | Map(lambda e: e[1]) | Map(self._drop_unwanted_fields)

        return events_to_write | "toBQ" >> WriteToBigQuery(
            table=self._choose_table,
            table_side_inputs=(table_names_dct,),
            create_disposition=BigQueryDisposition.CREATE_NEVER,
            insert_retry_strategy=RetryStrategy.RETRY_NEVER,
        )

您可以看到侧输入取自管道table_name 的另一个分支,该分支基本上是从事件中提取表名。然后,将其作为输入提供给WriteToBigQuery。不幸的是,这在负载下并不能真正起作用,侧输入没有更新,并且某些事件使用了错误的目的地。

在这种特定情况下,我还可以使用哪些其他方法?所有文档都使用静态示例,并没有真正涵盖这种动态方法。

我尝试的另一件事是编写一个使用 HTTP BigQuery 客户端并插入行的自定义 DoFn,这里的问题是管道的速度,因为每秒插入大约 6-7 个事件。

【问题讨论】:

    标签: python-3.x google-cloud-dataflow apache-beam


    【解决方案1】:

    我有一个类似的问题,我有一个解决办法。

    我看到你有create_disposition=BigQueryDisposition.CREATE_NEVER,所以在代码运行之前表列表是已知的。也许它笨拙,但它是众所周知的。我有一个DoFn 其中yeilds 很多TaggedOutputs 它的process 方法。然后我的管道看起来像:

    parser_outputs = ['my', 'list', 'of', 'tables']
    with beam.Pipeline(options=PipelineOptions(), argv=args) as p:
        pipe = (
            p
            | "Start" >> beam.Create(["example row"])
            | "Split"
            >> beam.ParDo(MySplitFn()).with_outputs(*parser_outputs)
        )
    
        for output in parser_outputs:
            pipe[output] | "write {}".format(output) >> beam.io.WriteToBigQuery(
                bigquery.TableReference(
                    projectId=options.projectId, datasetId=DATASET_ID, tableId=output
                ),
                schema=padl_shared.getBQSchema(parser.getSchemaForDataflow(rowTypeName=output)),
                write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            )
    
        p.run().wait_until_finish() 
    

    【讨论】:

    • 这是否可以很好地扩展?我的意思是,我有大约 23 个表应该在管道中写入。
    • 好的,我试过了,虽然它有效,但主要问题是它标记值的步骤非常慢,吞吐量约为 6-7 个元素/秒。不太明白为什么这一步这么慢,因为它只是对事件做了一个小的转换,然后用一个标签返回它。
    • 我的 DoFn 功能相当复杂。它包括来自导入模块的对象,导入对象的单元测试几乎有 1000 行长(对象本身只有 200 行加上一个配置文件)。我仍然在单个工作节点上获得 1000 多个元素/秒。如果你run your function locally in a test framework is it still slow?
    • 本地运行正常,而且我无法模拟生产环境的负载。另外,我注意到我收到了这个异常apache_beam.coders.coder_impl.IntervalWindowCoderImpl.estimate_size: TypeError: Cannot convert GlobalWindow to apache_beam.utils.windowed_value._IntervalWindowBase,它只发生在生产中,根本不会发生在本地。
    • 恐怕我不知道那个错误是关于什么的。我只将此模式用于批处理,而不是用于流式传输。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-05-16
    • 2020-06-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-25
    • 2020-10-29
    相关资源
    最近更新 更多