【问题标题】:How should a ZeroMQ worker safely "hang up"?ZeroMQ 工作者应该如何安全地“挂断”?
【发布时间】:2011-04-11 04:59:23
【问题描述】:

我本周开始使用 ZeroMQ,当使用请求-响应模式时,我不确定如何让工作人员安全地“挂断”并关闭他的套接字,而不会丢失消息并导致发送该消息的客户永远得不到回应。想象一个用 Python 编写的工人看起来像这样:

import zmq
c = zmq.Context()
s = c.socket(zmq.REP)
s.connect('tcp://127.0.0.1:9999')
while i in range(8):
    s.recv()
    s.send('reply')
s.close()

我一直在做实验,发现127.0.0.1:9999 的套接字类型为zmq.REQ 的客户提出公平排队请求可能会不幸让公平排队算法在worker 已经完成了最后一个send(),但在它运行以下close() 方法之前。在这种情况下,似乎请求被工作进程中的ØMQ堆栈接收和缓冲,然后当close()抛出与套接字关联的所有内容时,请求就会丢失。

工人如何“安全地”分离——有没有办法发出“我不再需要消息”的信号,然后 (a) 循环在信号传输期间到达的任何最终消息,(b) 生成他们的回复,然后 (c) 执行 close() 并保证没有消息被丢弃?

编辑:我想我想进入的原始状态是“半关闭”状态,在这种状态下无法接收到更多请求——并且发送者会知道——但是在哪里返回路径仍处于打开状态,因此我可以检查传入缓冲区中是否有最后到达的消息,如果缓冲区中有一条消息,则对其进行响应。

编辑:为了回答一个好问题,更正了描述,使等待消息的数量成为复数,因为可能有许多连接在等待回复。

【问题讨论】:

  • 我找到了这个线程,它是相关但不是真的:mail-archive.com/zeromq-dev@lists.zeromq.org/msg02338.html。您在这方面取得了进展吗?
  • 不,我还没有取得进一步的进展,ZeroMQ 的专业知识似乎比我预期的要薄一些。感谢您的链接 - 当消息出站时,它似乎涉及另一个方向的消息安全,但会带来许多相同的问题。谢谢!

标签: python concurrency message-queue rpc zeromq


【解决方案1】:

您似乎认为您正在尝试避免“简单”的竞争条件,例如 in

... = zmq_recv(fd);
do_something();
zmq_send(fd, answer);
/* Let's hope a new request does not arrive just now, please close it quickly! */
zmq_close(fd);

但我认为问题在于公平排队(循环)使事情变得更加困难:您的工作人员甚至可能已经有几个排队的请求。如果轮到它接收一个新请求,发件人不会等待您的工作人员空闲后再发送一个新请求,因此在您调用 zmq_send 时,其他请求可能已经在等待。

事实上,您可能选择了错误的数据方向。与其让请求池向您的工作人员发送请求(即使您不想接收新请求),您可能希望让您的工作人员从请求队列中获取新请求,处理它,然后发送答案。

当然,这意味着使用XREP/XREQ,但我认为这是值得的。

编辑:我写了some code实现the other direction来解释我的意思。

【讨论】:

  • 你说的很对;在一个比我的消息中最初描述的更复杂的示例中,可能有许多客户端,可能有许多等待消息。我已经更新了文本以反映这种可能性。谢谢!
  • 我用我写的一些新代码更新了我的答案,以展示如何处理您的问题。
【解决方案2】:

我认为问题在于您的消息传递架构是错误的。您的工作人员应该使用 REQ 套接字发送工作请求,这样只有一个工作在工作人员排队。然后,要确认工作完成,您可以使用另一个 REQ 请求,该请求兼作前一个作业的 ack 并请求一个新作业,或者您可以使用第二个控制套接字。

有些人使用 PUB/SUB 进行控制,这样每个工作人员都会发布确认信息,而主服务器会订阅它们。

您必须记住,使用 ZeroMQ 时有 0 个消息队列。一个都没有!只是消息在发送方或接收方中缓冲,具体取决于高水位标记和套接字类型等设置。如果您确实需要消息队列,那么您需要编写一个代理应用程序来处理它,或者简单地切换到所有通信都通过第三方代理进行的 AMQP。

【讨论】:

  • 很好的解释——感谢您对我为什么试图以错误的方式做事的一些哲学。我在想我的工人应该被工作,然后会有一些神奇的方式让它说“停止”——但实际上解决方案是让工人每次都开始对话,这使它处于控制之中,而不是让发送者处于控制之中。谢谢!
【解决方案3】:

我也一直在考虑这个问题。您可能希望实现一条 CLOSE 消息,通知客户工人即将离开。然后,您可以让工作人员在关闭之前耗尽一段时间。当然不理想,但可能可行。

【讨论】:

  • 是的,这是可能的。但我希望使用一个 REQ 客户,该客户甚至不知道我设置了多少台服务器来对他的请求进行负载平衡。我希望避免必须迁移到 XREQ 并实施我自己的订阅和取消订阅操作!
  • 不是一个很好的解决方案,但目前似乎不存在任何真正的解决方案:S.
【解决方案4】:

在尽可能快地向工作人员发送请求与在工作崩溃或死亡的情况下获得可靠性之间存在利益冲突。 ZeroMQ 指南有一整节解释了这个可靠性问题的不同答案。读一下,会有很大帮助的。

tl;dr 工作人员可能/将会崩溃,并且客户端需要重新发送功能。该指南提供了多种语言的可重用代码。

【讨论】:

  • 这一定是指南中比我写这个问题时更新的部分!我会看看他们现在的建议是什么。谢谢!
【解决方案5】:

最简单的解决方案不是让客户在等待回复时超时,然后在没有收到回复时重试吗?

【讨论】:

  • 如何为 ØMQ recv() 函数指定超时时间?我认为它永远阻塞,没有任何逃脱的机会。
  • @Brandon:你可以用zmq_poll() 来做到这一点。不过,不必编写这种超时和恢复代码是 ZMQ 的优势之一。 :(
  • 重试是个问题,因为它们会在某个特定工作比正常时间更长的情况下重复工作。请注意,在这些情况下,多次提交的是较长的作业,从而增加了提供给服务器的负载。这里的问题是,不像 TCP,它采用不可靠的网络,让你几乎忘记数据包是否丢失的问题,ØMQ 采用 TCP,让你再次开始担心丢失的消息。在这个 REQ-REP 用例中,它似乎是应用程序复杂性的一大步。
【解决方案6】:

尝试在调用关闭之前睡觉。这在 2.1 中已修复,但在 2.0 中尚未修复。

【讨论】:

  • 睡眠同步让我害怕。你知道这在 2.1 中是如何解决的吗?他们是否为半关闭套接字添加了一个选项,或者关闭是否告诉发送者他们需要在其他地方重新传输的未处理消息?
  • Trey,我看到在 2.1 中 close() 不会破坏排队的传出消息。但是我看不到传入消息,至少在我正在查看的小摘要页面上看不到。您能否将我们指向正确位置的更改日志或文档?
  • 对不起,伙计们,我有同样的信息。我不是 0mq 的贡献者。据我所知,2.1 只是在关闭时刷新而不是立即关闭。
猜你喜欢
  • 2020-10-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-02-21
  • 2021-08-01
  • 2015-07-15
  • 2014-09-12
  • 2011-06-23
相关资源
最近更新 更多