【发布时间】: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?您如何发布带有时间戳标签的消息?水印卡在什么地方?
-
感谢您的回复本。我用您的问题的答案更新了我的初始帖子。