Beam / Dataflow设计模式可根据数据库查询来丰富文档

sev*_*ian 4 google-cloud-dataflow apache-beam

评估数据流,并试图找出是否/如何执行以下操作。

我很抱歉,如果上面的内容微不足道,请在我们决定使用Beam或Spark等其他东西之前,先将头绕在Dataflow上。

一般用例用于机器学习:

  • 摄取单独处理的文档。

  • 除了易于编写的转换,我们还要基于对数据库(主要是键值存储)的查询来丰富每个文档。

  • 一个简单的例子是地名词典:将文本分解为ngram,然后检查这些ngram是否驻留在某个数据库中,并记录(在原始文档的转换版本中)给定短语映射到的实体标识符。

如何有效地做到这一点?

天真(尽管可能对序列化要求很棘手?):

每个文档都可以简单地单独查询数据库(类似于通过Google DataFlow Transformer查询关系数据库),但是,鉴于大多数文件都是简单的键值存储,因此似乎应该有一种更有效的方法(鉴于数据库查询延迟的实际问题)。

场景1:已改进?:

当前的稻草人将表存储在Bigquery中,将它们拉下来(https://github.com/apache/beam/blob/master/sdks/python/apache_beam/io/gcp/bigquery.py),然后将它们用作侧面输入,用作每个文档函数中的键值查找。

键值表的范围通常从很小到不大(100 MB,甚至是低GB)。 多个CoGroupByKey具有相同的apache光束(“侧面输入可以任意大-没有限制;我们已经看到使用1 + TB大小的侧面输入成功地运行了管道”)表明这是合理的,至少从POV大小而言。

1)这有意义吗?这是此方案的“正确”设计模式吗?

2)如果这是一个好的设计模式...我该如何实际实现呢?

https://github.com/apache/beam/blob/master/sdks/python/apache_beam/io/gcp/bigquery.py#L53显示将结果作为AsList馈入文档功能。

i)大概,AsDict在这里更适合上述用例?因此,我可能需要首先在Bigquery输出上运行一些转换,以将其分成键,值元组;并确保密钥是唯一的;然后将其用作侧面输入。

ii)然后我需要在函数中使用侧面输入。

我不清楚的是:

  • 对于这两个方面,如何处理Bigquery拉动产生的输出对我来说还是很模糊的。我将如何完成(i)(假设有必要)?意思是,数据格式是什么样的(原始字节?字符串?我可以研究一个很好的例子吗?)

  • 同样,如果AsDict是将其传递给func的正确方法,我是否可以引用通常在python中使用的dict这样的内容?例如,side_input.get('blah')吗?

场景2:进一步改善了吗?(针对特定情况):

  • 上面的场景(如果可以实现的话)肯定确实是高级连续远程调用(考虑到简单的键值查找),这对于我们的某些场景非常有帮助。但是,如果我采用类似地名词典查找的方案(如上)……是否有更优化的解决方案?

像这样,对于每个文档,将我们所有的ngram都写为键,将值作为基础索引(文档中的docid + indices),然后在这些ngram和gazeteer中的短语之间进行某种连接...然后进行另一组转换以恢复原始文档(现在带有新注释)。

即,让Beam直接处理所有联接/查找吗?

从理论上讲,Beam比在每个文档中循环遍历所有ngram并检查ngram是否在side_input中的速度要快得多。

其他关键问题:

3)如果这是做事的好方法,是否有任何技巧可以使它在流媒体情况下正常工作?其他地方的文字表明,在批处理方案之外,侧面输入缓存的工作效果更差。目前,我们专注于批量处理,但是流式传输在提供实时预测方面将变得至关重要。

4)是否有任何与Beam有关的原因可以为上述任何原因选择Java> Python?我们已经有大量现有的Python代码移至Dataflow,因此会非常喜欢Python ...但不确定上面的Python是否存在任何隐藏问题(例如,我注意到Python不支持某些功能或I / O)。

编辑:稻草人?对于示例ngram查找场景(应大力推广到常规K:V查找)

  • 短语=从bigquery获取
  • Docs(由docid编制索引)(例如,来自文本或protobufs的直接输入)
  • 转换:词组->(短语,实体)元组
  • 转换:docs-> ngrams(短语,docid,坐标[在文档中])
  • CoGroupByKey key =短语:(短语,实体,docid,coords)
  • CoGroupByKey key = docid,group((短语,实体,docid,coords),Docs)
  • 然后,我们可以使用(短语,实体,docid,coords)和每个文档的集合来迭代地完成每个文档

Pab*_*blo 5

关于管道的方案:

  1. 天真的场景

没错,按元素查询数据库是不正确的。

如果您的键值存储能够通过重用打开的连接来支持低延迟查找,则可以定义一个全局连接,该全局连接为每个工作线程初始化一次,而不是为每个捆绑包初始化一次。您的kv商店支持对现有连接的有效查找,这应该是可以接受的。

  1. 改进方案

如果那不可行,那么BQ是保存和提取数据的好方法。

您绝对可以使用AsDict侧面输入,只需转到side_input[my_key]或即可side_input.get(my_key)。

您的管道可能看起来像这样:

kv_query = "SELECT key, value FROM my:table.name"
p = beam.Pipeline()
documents_pcoll = p | ReadDocuments()
additional_data_pcoll = (p 
                   | beam.io.BigQuerySource(query=kv_query)
                   # Make row a key-value tuple.
                   | 'format bq' >> beam.Map(lambda row: (row['key'], row['value'])))

enriched_docs = (documents_pcoll 
                 | 'join' >> beam.Map(lambda doc, query: enrich_doc(doc, query[doc['key']]), 
                                      query=AsDict(additional_data_pcoll)))
Run Code Online (Sandbox Code Playgroud)

不幸的是,这有一个缺点,那就是Python当前不支持任意大的侧面输入(它当前将所有KV加载到一个Python字典中)。如果您的侧面输入数据很大,那么您将避免使用此选项。

注意将来这会改变,但是我们不能确定ATM。

  1. 进一步改善

连接两个数据集的另一种方法是使用CoGroupByKey。文档和KV附加数据的加载不应更改,但是在加入时,您将执行以下操作:

# Turn the documents into key-value tuples as well[
documents_kv_pcoll = (documents_pcoll 
                      | 'format docs' >> beam.Map(lambda doc: (doc['key'], doc)))
enriched_docs = ({'docs': documents_kv_pcoll, 'additional_data': additional_data_pcoll}
                 | beam.CoGroupByKey()
                 | 'enrich' >> beam.Map(lambda x: enrich_doc(x['docs'][0], x['additional_data'][0]))
Run Code Online (Sandbox Code Playgroud)

CoGroupByKey 将允许您在任一侧使用任意大的集合。

回答你的问题

  1. 您可以在本手册中看到一个使用BigQuery作为侧面输入的示例。如您所见,数据被解析(我相信它来自其原始数据类型,但可能以字符串/ Unicode形式出现)。查看文档(或随时询问)是否需要更多信息。

  2. 目前,Python流处于Alpha中,它不支持侧面输入。但它确实支持随机播放功能,例如CoGroupByKey。您使用的管道CoGroupByKey在流式传输中应该工作良好。

  3. 之所以选择Java而不是Python,是因为所有这些功能都可以在Java中使用(大小不受限制的侧面输入,流式侧面输入)。但似乎对于您的用例,Python可能具有您所需要的全部。

注意:这些代码段是近似的,但是您应该可以使用进行调试DirectRunner。

如果您认为有帮助,请随时进行澄清,或询问其他方面。