【发布时间】: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