【发布时间】:2019-07-14 15:18:22
【问题描述】:
我的流束/数据流管道正在通过 pub/sub 从另一个服务一个接一个地接收基于事件的数据。为了确保对上游数据结构进行更改的任何人都不会破坏管道,我在每个元素上运行以下代码:
class CreateLoadsTableRow(beam.DoFn):
def process(self, element):
row = {
'event_id': element.get('load_id'),
'domain': element.get('url'),
'user_data': {
'event_id': element.get('events'),
}
# Loads more keys below
}
yield row
我担心这会非常昂贵 - 有没有更有效的方法来实现这一点?
或者有没有更好的模式?
【问题讨论】:
标签: python google-bigquery google-cloud-dataflow apache-beam