【问题标题】:How to pause and resume consumption gracefully in rabbitmq, pika python如何在rabbitmq、pika python中优雅地暂停和恢复消费
【发布时间】:2013-12-18 06:30:56
【问题描述】:

我用basic_consume()接收消息,basic_cancel取消消费,但是有问题。

这里是 pika.channel 的代码

 def basic_consume(self, consumer_callback, queue='', no_ack=False,
                      exclusive=False, consumer_tag=None):
        """Sends the AMQP command Basic.Consume to the broker and binds messages
        for the consumer_tag to the consumer callback. If you do not pass in
        a consumer_tag, one will be automatically generated for you. Returns
        the consumer tag.

        For more information on basic_consume, see:
        http://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

        :param method consumer_callback: The method to callback when consuming
        :param queue: The queue to consume from
        :type queue: str or unicode
        :param bool no_ack: Tell the broker to not expect a response
        :param bool exclusive: Don't allow other consumers on the queue
        :param consumer_tag: Specify your own consumer tag
        :type consumer_tag: str or unicode
        :rtype: str

        """
        self._validate_channel_and_callback(consumer_callback)

        # If a consumer tag was not passed, create one
        consumer_tag = consumer_tag or 'ctag%i.%s' % (self.channel_number,
                                                      uuid.uuid4().get_hex())

        if consumer_tag in self._consumers or consumer_tag in self._cancelled:
            raise exceptions.DuplicateConsumerTag(consumer_tag)

        self._consumers[consumer_tag] = consumer_callback
        self._pending[consumer_tag] = list()
        self._rpc(spec.Basic.Consume(queue=queue,
                                     consumer_tag=consumer_tag,
                                     no_ack=no_ack,
                                     exclusive=exclusive),
                           self._on_eventok,
                           [(spec.Basic.ConsumeOk,
                             {'consumer_tag': consumer_tag})])

        return consumer_tag

def basic_cancel(self, callback=None, consumer_tag='', nowait=False):
        """This method cancels a consumer. This does not affect already
        delivered messages, but it does mean the server will not send any more
        messages for that consumer. The client may receive an arbitrary number
        of messages in between sending the cancel method and receiving the
        cancel-ok reply. It may also be sent from the server to the client in
        the event of the consumer being unexpectedly cancelled (i.e. cancelled
        for any reason other than the server receiving the corresponding
        basic.cancel from the client). This allows clients to be notified of
        the loss of consumers due to events such as queue deletion.

        :param method callback: Method to call for a Basic.CancelOk response
        :param str consumer_tag: Identifier for the consumer
        :param bool nowait: Do not expect a Basic.CancelOk response
        :raises: ValueError

        """
        self._validate_channel_and_callback(callback)
        if consumer_tag not in self.consumer_tags:
            return
        if callback:
            if nowait is True:
                raise ValueError('Can not pass a callback if nowait is True')
            self.callbacks.add(self.channel_number,
                               spec.Basic.CancelOk,
                               callback)
        self._cancelled.append(consumer_tag)
        self._rpc(spec.Basic.Cancel(consumer_tag=consumer_tag,
                                    nowait=nowait),
                  self._on_cancelok,
                  [(spec.Basic.CancelOk,
                    {'consumer_tag': consumer_tag})] if nowait is False else [])

正如您所看到的,每次我取消消费时都会将 consumer_tag 添加到 _canceled 列表中。如果我再次在 basic_consume 中使用这个标签,就会引发 duplicateConsumer 异常。 好吧,我每次都可以使用一个新的 consumer_tag,但实际上我不是。因为迟早生成的标签会与之前的一些标签完全匹配。

如何在 pika 中优雅地暂停和恢复消费?

【问题讨论】:

    标签: python rabbitmq pika


    【解决方案1】:

    您定义自己的consumer_tags 有什么原因吗?你可以传递一个空字符串,让 RabbitMQ 为你生成消费者标签。来自basic.consume的回复,也就是basic.consume-ok会返回生成的consumer_tag,以后可以用它来停止消费。

    见:http://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume-ok

    【讨论】:

    • 从 pika 模块的附加代码可以看出,即使在客户端应用程序没有指定消费者标签的情况下,那么 pika 模块也会生成一个并将其包含在消费请求中.
    【解决方案2】:

    看起来 Pika 做得比它应该做的要多 - 如果没有提供消费者标签,它不需要创建消费者标签(服务器会),它也不需要注意重复的消费者标签(恢复服务器支持相同的标签)。

    所以我不确定如何使用 Pika 执行此操作 - 我想提交一个错误。

    【讨论】:

      猜你喜欢
      • 2018-06-10
      • 2022-01-23
      • 1970-01-01
      • 2017-01-27
      • 1970-01-01
      • 2019-09-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多