【发布时间】:2017-01-27 18:55:32
【问题描述】:
根据堆栈驱动程序图表,我们开始注意到某个主题/订阅的“未确认消息”数量不时增加。
症状
我不知道我们可以在多大程度上信任 stackdriver 图表,但我已经检查过:
- 拉取操作数与发布操作数一样多
- 出现问题时,ack 操作计数低于 pull 操作计数
此外,根据我们的日志,我能够看到 pubsub 实际上多次发送相同的消息,这也证实了“pull”成功但“ack”可能不成功。
所以,我认为我们可以假设我们的系统会迅速拉取,但从 GCP 的角度来看并不能很好地 ACK。
我检查了没有按时发送 ACK 的可能性,但我认为不是这样,如下面的流程所示。
在有问题的订阅中,消息被累积了几个小时。对我们来说,这是一个严重的问题。
实现细节
我们出于某种原因使用 pull 方法,并且我们不愿意切换到 push 方法,除非有充分的理由。对于每个订阅,我们有一个消息泵 goroutine,这个 goroutine 为每条拉取的消息生成一个 worker。更具体地说,
// in a dedicated message-pumping goroutine
sub, _ := CreateSubscription(..., 0 /* ack-deadline */, )
iter, _ := sub.Pull(...)
for {
// omitted: wait if we have too many workers
msg, _ := iter.Next()
go func(msg Message) {
// omitted: handle the message and measure the latency; it turned out the latency is almost within 1 second
msg.Done(true)
}(msg)
}
对于负载平衡,订阅也会被同一集群中的其他 Pod 拉取。因此,对于一个订阅(如在 Google Pubsub 主题/订阅中),我们有多个订阅对象(如在 Go 绑定的订阅结构中),每个订阅对象仅在一个 pod 中使用。并且,每个订阅对象都会创建一个迭代器。我相信这个设置没有错,但如果我错了,请纠正我。
正如这段代码所示,我们执行 ACK。 (我们的服务器不会恐慌;因此没有绕过 msg.Done() 的途径。)
尝试
奇怪的是,有问题的订阅并不忙。对于在同一个 pod 中接收更多消息的另一个订阅,我们通常不会有任何问题。因此,我开始怀疑 pull 操作的 max-prefetch 选项是否会影响。似乎解决了一段时间的问题,但问题再次出现。
按照 Google 支持的建议,我还增加了 pod 的数量,这有效地增加了工作人员的数量。这没有多大帮助。由于我们没有向有问题的问题发布很多消息(大约 1 条消息/秒),而且我们有很多(可能太多)工作人员,我认为我们的服务器不会超载。
有人能解释一下吗?
【问题讨论】:
-
我在上周开始使用 Node.js 库时遇到了同样的问题。 ack 不起作用,然后我必须等待消息重新传递