【问题标题】:Watermark getting stuck水印卡住了
【发布时间】:2017-10-02 09:43:25
【问题描述】:

我正在通过 pub/sub 将数据摄取到以无限制模式运行的数据流管道。这些数据基本上是与从跟踪设备捕获的时间戳的坐标。这些消息分批到达,每批可能是 1..n 条消息。在一段时间内可能没有消息到达,稍后可能会重新发送(或不重新发送)。我们使用每个坐标的时间戳(以 UTC 为单位)作为 pub-sub 消息的属性。并通过时间戳标签读取管道:

pipeline.apply(PubsubIO.Read.topic("new").timestampLabel("timestamp")

坐标和延迟的示例如下:

36 points wait 0:02:24
36 points wait 0:02:55
18 points wait 0:00:45
05 points wait 0:00:01
36 points wait 0:00:33
36 points wait 0:00:43
36 points wait 0:00:34

消息可能如下所示:

2013-07-07 09:34:11;47.798766;13.050133

在第一批之后,水印是空的,在第二批之后,我可以在管道诊断中看到一个水印,只是它没有得到更新,尽管有新消息到达。同样根据堆栈驱动程序日志记录,PubSub 没有未传递或未确认的消息。

随着具有新事件时间的消息到达,水印不应该向前移动吗?

根据What is the watermark heuristic for PubsubIO running on GCD?,WaterMark 也应该每 2 分钟向前移动一次,而不是?

[..] 如果我们没有看到更多关于订阅的数据 超过两分钟(并且没有积压),我们将水印推进到 接近实时。 [..]

更新以解决 Bens 的问题:

是否有我们可以查看的作业 ID?

是的,我刚刚在 09:52 CET 即 07:52 UTC 重新启动了整个设置,作业 ID 为 2017-05-05_00_49_11-11176509843641901704。

您使用的是什么版本的 SDK?

1.9.0

您如何发布带有时间戳标签的消息?

我们使用 python 脚本发布使用 pub sub sdk 的数据。 来自那里的消息可能如下所示:

{'data': {timestamp;lat;long;ele}, 'timestamp': '2017-05-05T07:45:51Z'}

我们在数据流中为时间戳标签使用时间戳属性。

水印卡在什么地方?

对于这项工作,水印现在停留在 09:57:35(我在 10:10 左右发布),尽管发送了新数据,例如在

10:05:14
10:05:43
10:06:30

我还可以看到,我们可能会在延迟超过 10 秒的情况下将数据发布到 pub sub,例如我们在 10:07:47 发布最高时间戳为 10:07:26 的数据。

几个小时后,水印赶上来了,但我不明白为什么它会延迟/一开始没有移动。

【问题讨论】:

  • 是否有工作 ID 可供我们查看?您使用的是什么版本的 SDK?您如何发布带有时间戳标签的消息?水印卡在什么地方?
  • 感谢您的回复本。我用您的问题的答案更新了我的初始帖子。

标签: google-cloud-dataflow


【解决方案1】:

这是 PubSub 水印跟踪逻辑中的一个边缘案例,它有两种变通方法(见下文)。本质上,如果 2 分钟内没有 no 输入,则水印将前进到当前时间。但是,如果数据的到达速度快于每 2 分钟一次,但 QPS 仍然非常低,那么就没有足够的数据来使估计的水印保持最新。

正如我所提到的,有几种解决方法:

  1. 如果您处理更多数据,问题自然会得到解决。
  2. 或者,如果您注入额外的消息(例如每秒 2 条),它将为水印提供足够的数据以更快地推进。这些只需要有时间戳,并且可以立即从管道中过滤掉。

【讨论】:

  • 谢谢。我在尝试管道时遇到了这个问题,不得不推送虚拟消息以使水印保持最新(当然我会收到更多消息)。应该更新数据流文档以包含除 DF 工程师之外没有人知道的“边缘案例”,尤其是当它们违反最小惊讶原则时。
【解决方案2】:

作为记录,关于前面提到的 直接运行器 上下文中的边缘情况,另一件要记住的事情是运行器的并行性。具有更高的并行度(尤其是在多核机器上是默认设置)似乎需要更多的数据。就我而言,设置--targetParallelism=1 有所帮助。基本上在没有任何其他干预的情况下将卡住的管道转变为正常工作的管道。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-11-09
    • 2012-06-23
    • 2021-10-08
    • 2016-11-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多