【发布时间】:2017-02-03 15:35:18
【问题描述】:
我正在使用 Spotify Scio 读取从 Stackdriver 导出到 Google Cloud Storage 的日志。它们是 JSON 文件,其中每一行都是一个条目。查看工作日志,文件似乎被分成块,然后以任何顺序读取。在这种情况下,我已经将我的工作限制为 1 名工人。有没有办法强制按顺序读取和处理这些块?
举个例子(textFile基本上就是一个TextIO.Read):
val sc = ScioContext(myOptions)
sc.textFile(myFile).map(line => logger.info(line))
会根据工作日志产生与此类似的输出:
line 5
line 6
line 7
line 8
<Some other work>
line 1
line 2
line 3
line 4
<Some other work>
line 9
line 10
line 11
line 12
我想知道是否有办法强制它按顺序读取第 1-12 行。我发现压缩文件并使用指定的 CompressionType 读取它是一种解决方法,但我想知道是否有任何方法可以做到这一点,而不涉及压缩或更改原始文件。
【问题讨论】:
-
我最近遇到了类似的问题,反馈基本上是“否”。似乎即使您在本地运行,Dataflow 仍然以随机顺序读取。我为此实施的一种解决方法不是很好,它是在 Pub/Sub 中按顺序读取文件并将消息发送到 Dataflow,Dataflow 订阅 PubSub 主题,而不是读取文件。当 DataFlow 完成每条消息时,它会发回消息说它已经完成,因此 PubSub 发送下一条消息。这有点过火了,所以很高兴听到更好/内置的选项......
-
很不幸,我一直在考虑做类似的事情,但我可能只是预压缩所有内容,因为它至少看起来可靠。如果我理解正确,则在压缩它们并且逻辑按顺序进行时不会发生拆分。我同意应该有一个更简单的方法,谢谢!
-
您能详细说明您的用例吗? Dataflow 用于数据并行处理,听起来您正在寻找串行工具。
-
+1 回答 Sam 的问题。请不要依赖观察到的对 zip 文件的有序处理(或任何其他观察到的行为,Beam 编程模型未明确保证) - 这可能随时更改,恕不另行通知;即使您使用相同的代码和相同的 SDK 版本,它也会发生变化,因为我们在后端实施了新的优化。
-
我正在尝试使用导出的 Stackdriver 日志,通过我在 Dataflow 中的转换逻辑运行它们,然后将它们转发给其他 GC 服务。感谢 Sam 和 jkff 的输入,我想我可以稍微改变一下设计,让数据流之外的顺序部分。
标签: google-cloud-platform google-cloud-dataflow spotify-scio