【问题标题】:Dataflow batch vs streaming: window size larger than batch size数据流批处理与流式处理:窗口大小大于批处理大小
【发布时间】:2018-02-06 03:19:28
【问题描述】:

假设我们有带有时间戳的日志数据,这些数据可以流式传输到 BigQuery 或作为文件存储在 Google Storage 中,但不能直接流式传输到 Dataflow 支持的无限收集源类型。

我们想根据时间戳来分析这些数据,无论是相对的还是绝对的,例如“最近 1 小时内点击了多少?”和“2018 年 2 月 5 日下午 3 点到 4 点之间有多少点击量”?

阅读了有关窗口和触发器的文档后,尚不清楚如果我们想要一个大窗口,我们将如何以 Dataflow 支持的方式将传入数据分成批次 - 可能我们希望在最后一天进行聚合、30 天、3 个月等。

例如,如果我们的批处理源是 BigQuery 查询,每 5 分钟运行一次,对于最后 5 分钟的数据,Dataflow 是否会在作业运行之间保持窗口打开,即使数据以 5 分钟的块到达?

同样,如果日志文件每 5 分钟轮换一次,并且我们在将新文件保存到存储桶时启动 Dataflow,则同样的问题适用 - 作业是否已停止和启动,并且先前作业的所有知识都已丢弃,或者大窗口(例如最多一个月)是否仍然为新活动打开?

我们如何在不干扰现有状态的情况下更改/修改此管道?

如果这些是基本问题,我们深表歉意,即使是指向某些文档的链接也将不胜感激。

【问题讨论】:

    标签: google-cloud-dataflow apache-beam


    【解决方案1】:

    听起来您想要对数据进行任意交互式聚合查询。 Beam / Dataflow 本身并不适合这种情况,但是 Dataflow 最常见的用例之一是将数据提取到 BigQuery(例如来自 GCS 文件或来自 Pubsub)中,非常适合。

    关于您的问题的更多信息:

    目前尚不清楚我们如何将传入的数据分成批次

    Beam 中的窗口化只是一种在时间维度上指定聚合范围的方法。例如。如果您每 5 分钟使用大小为 15 分钟的滑动窗口,则事件时间时间戳为 14:03 的记录计入三个窗口中的聚合:13:50..14:05、13:55..14: 10点,14:00..14:15。

    所以:与您在按键分组时不需要将传入数据划分为“键”的方式相同(数据处理框架为您逐键执行分组),您不需要将其划分为windows 也可以(框架在每个聚合操作中隐式执行逐个窗口分组)。

    Dataflow 会在作业运行之间保持窗口打开

    我希望前一点已解决此问题,但要澄清更多:不。停止 Dataflow 作业会丢弃其所有状态。但是,您可以使用新代码“更新”作业(例如,如果您修复了错误或添加了额外的处理步骤) - 在这种情况下,状态不会被丢弃,但我认为这不是您要问的。

    如果日志文件每 5 分钟轮换一次,并且我们在保存新文件时启动 Dataflow

    听起来您想连续提取数据。做到这一点的方法是编写一个连续运行的流式管道来连续地摄取数据,而不是在每次新数据到达时启动一个新的管道。在文件到达存储桶的情况下,如果您正在阅读文本文件,则可以使用 TextIO.read().watchForNewFiles();如果您正在阅读其他类型的文件,则可以使用它的各种类似物(最常见的是 FileIO.matchAll().continuously())。

    【讨论】:

    • “但我认为这不是你要问的” -> 谢谢,这就是我要问的。我一定会跟进这些功能。我们确实想连续摄取,所以它们应该是合适的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-28
    • 1970-01-01
    • 1970-01-01
    • 2014-01-19
    • 2017-02-01
    • 1970-01-01
    相关资源
    最近更新 更多