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查找)
关于管道的方案:
没错,按元素查询数据库是不正确的。
如果您的键值存储能够通过重用打开的连接来支持低延迟查找,则可以定义一个全局连接,该全局连接为每个工作线程初始化一次,而不是为每个捆绑包初始化一次。您的kv商店支持对现有连接的有效查找,这应该是可以接受的。
如果那不可行,那么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。
连接两个数据集的另一种方法是使用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 将允许您在任一侧使用任意大的集合。
您可以在本手册中看到一个使用BigQuery作为侧面输入的示例。如您所见,数据被解析(我相信它来自其原始数据类型,但可能以字符串/ Unicode形式出现)。查看文档(或随时询问)是否需要更多信息。
目前,Python流处于Alpha中,它不支持侧面输入。但它确实支持随机播放功能,例如CoGroupByKey。您使用的管道CoGroupByKey在流式传输中应该工作良好。
之所以选择Java而不是Python,是因为所有这些功能都可以在Java中使用(大小不受限制的侧面输入,流式侧面输入)。但似乎对于您的用例,Python可能具有您所需要的全部。
注意:这些代码段是近似的,但是您应该可以使用进行调试DirectRunner。
如果您认为有帮助,请随时进行澄清,或询问其他方面。
| 归档时间: |
|
| 查看次数: |
1023 次 |
| 最近记录: |