【问题标题】:How to get all messages in Amazon SQS queue using boto library in Python?如何使用 Python 中的 boto 库获取 Amazon SQS 队列中的所有消息?
【发布时间】:2012-04-28 04:19:08
【问题描述】:

我正在开发一个应用程序,其工作流是通过使用 boto 在 SQS 中传递消息来管理的。

我的 SQS 队列正在逐渐增长,我无法检查它应该包含多少元素。

现在我有一个定期轮询队列的守护进程,并检查我是否有一组固定大小的元素。例如,考虑以下“队列”:

q = ["msg1_comp1", "msg2_comp1", "msg1_comp2", "msg3_comp1", "msg2_comp2"]

现在我想在某个时间点检查队列中是否有“msg1_comp1”、“msg2_comp1”和“msg3_comp1”,但我不知道队列的大小。

查看 API 后,您似乎只能获取 1 个元素,或队列中的固定数量的元素,但不是全部:

>>> rs = q.get_messages()
>>> len(rs)
1
>>> rs = q.get_messages(10)
>>> len(rs)
10

答案中提出的建议是,例如在循环中获取 10 条消息,直到我什么也得不到,但是 SQS 中的消息有可见性超时,这意味着如果我从队列中轮询元素,它们将不会真的被移除了,它们只会在短时间内不可见。

有没有一种简单的方法来获取队列中的所有消息,而不知道有多少?

【问题讨论】:

    标签: python boto amazon-sqs


    【解决方案1】:

    在while循环中拨打q.get_messages(n):

    all_messages=[]
    rs=q.get_messages(10)
    while len(rs)>0:
        all_messages.extend(rs)
        rs=q.get_messages(10)
    

    另外,dump won't support more than 10 messages 要么:

    def dump(self, file_name, page_size=10, vtimeout=10, sep='\n'):
        """Utility function to dump the messages in a queue to a file
        NOTE: Page size must be < 10 else SQS errors"""
    

    【讨论】:

    • 我真的做不到,因为 SQS 中的消息有可见性超时,所以如果我先收到 10 条消息,然后循环几次,下次我可能会收到相同的 10 条消息因为超时已经过去。我正在考虑使用dump(),但之后我必须阅读文件,这似乎很愚蠢,我错过了什么吗? (我可以将 visibility_timeout 设置为很长的时间,但这看起来很难看)。
    • @linker - 你说你需要检查'n'特定消息。这是否意味着您正在比较每条消息的某些匹配标准?
    • @linker - 根据参考资料,可见性超时最长可达 12 小时。除非您开始大规模的 EC2 工作,否则我猜这会满足您的需求吗? docs.amazonwebservices.com/AWSSimpleQueueService/2011-10-01/…
    • @linker - 顺便说一句,消息的数量应该是 1 到 10。如果你使用其他东西,SQS 服务应该返回一个 ReadCountOutOfRange 错误。
    • 在这种情况下依赖时间实际上有点棘手,因为我得到了一些我无法控制的事件(来自外部来源,如果他们有问题,例如停止发送可能造成问题的数据几个小时)。我不明白为什么有转储方法,但没有 get_all_messages,对我来说毫无意义。
    【解决方案2】:

    我一直在使用 AWS SQS 队列来提供即时通知,因此我需要实时处理所有消息。以下代码将帮助您有效地将(所有)消息出队并在删除时处理任何错误。

    注意:要从队列中删除消息,您需要删除它们。我正在使用更新的 boto3 AWS python SDK、json 库和以下默认值:

    import boto3
    import json
    
    region_name = 'us-east-1'
    queue_name = 'example-queue-12345'
    max_queue_messages = 10
    message_bodies = []
    aws_access_key_id = '<YOUR AWS ACCESS KEY ID>'
    aws_secret_access_key = '<YOUR AWS SECRET ACCESS KEY>'
    sqs = boto3.resource('sqs', region_name=region_name,
            aws_access_key_id=aws_access_key_id,
            aws_secret_access_key=aws_secret_access_key)
    queue = sqs.get_queue_by_name(QueueName=queue_name)
    while True:
        messages_to_delete = []
        for message in queue.receive_messages(
                MaxNumberOfMessages=max_queue_messages):
            # process message body
            body = json.loads(message.body)
            message_bodies.append(body)
            # add message to delete
            messages_to_delete.append({
                'Id': message.message_id,
                'ReceiptHandle': message.receipt_handle
            })
    
        # if you don't receive any notifications the
        # messages_to_delete list will be empty
        if len(messages_to_delete) == 0:
            break
        # delete messages to remove them from SQS queue
        # handle any errors
        else:
            delete_response = queue.delete_messages(
                    Entries=messages_to_delete)
    

    【讨论】:

    • 对 v2 Boto 软件包的改编以从 Boto3“反向移植”delete_messages 函数是 here。内置的 Boto(2) delete_message_batch 有 10 条消息的限制,并且需要完整的 Message 类对象,而不仅仅是对象中的 ID 和 ReceiptHandles。
    • 谢谢你的回答,对解决这个问题很有帮助stackoverflow.com/questions/62681836/…
    【解决方案3】:

    我的理解是,SQS 服务的分布式特性几乎使您的设计无法实现。每次调用 get_messages 时,您都在与一组不同的服务器通信,这些服务器将包含您的一些但不是全部消息。因此,不可能“不时签入”来设置特定消息组是否已准备好,然后只接受这些消息。

    您需要做的是不断地轮询,在所有消息到达时获取它们,并将它们本地存储在您自己的数据结构中。每次成功获取后,您可以检查您的数据结构以查看是否已收集到完整的消息集。

    请记住,消息将无序到达,并且某些消息将传递两次,因为删除必须传播到所有 SQS 服务器,但随后的 get请求有时会胜过删除消息。

    【讨论】:

      【解决方案4】:

      我在一个 cronjob 中执行这个

      from django.core.mail import EmailMessage
      from django.conf import settings
      import boto3
      import json
      
      sqs = boto3.resource('sqs', aws_access_key_id=settings.AWS_ACCESS_KEY_ID,
               aws_secret_access_key=settings.AWS_SECRET_ACCESS_KEY,
               region_name=settings.AWS_REGION)
      
      queue = sqs.get_queue_by_name(QueueName='email')
      messages = queue.receive_messages(MaxNumberOfMessages=10, WaitTimeSeconds=1)
      
      while len(messages) > 0:
          for message in messages:
              mail_body = json.loads(message.body)
              print("E-mail sent to: %s" % mail_body['to'])
              email = EmailMessage(mail_body['subject'], mail_body['message'], to=[mail_body['to']])
              email.send()
              message.delete()
      
          messages = queue.receive_messages(MaxNumberOfMessages=10, WaitTimeSeconds=1)
      

      【讨论】:

        【解决方案5】:

        注意:这不是直接回答问题。 相反,它是对@TimothyLiu's answer 的扩充,假设最终用户使用的是Boto 包(又名Boto2)而不是Boto3。此代码是 his answer 中提到的 delete_messages 调用的“Boto-2 化”


        Boto(2) 调用 delete_message_batch(messages_to_delete) 其中 messages_to_delete 是一个 dict 对象,其 key:value 对应于 id:receipt_handle 对返回

        AttributeError: 'dict' 对象没有属性 'id'。

        似乎delete_message_batch 需要一个Message 类对象;如果您一次删除超过 10 条“消息”,则复制 Boto source for delete_message_batch 并允许它使用非 Message 对象(ala boto3)也会失败。所以,我不得不使用以下解决方法。

        从here打印代码

        from __future__ import print_function
        import sys
        from itertools import islice
        
        def eprint(*args, **kwargs):
            print(*args, file=sys.stderr, **kwargs)
        
        @static_vars(counter=0)
        def take(n, iterable, reset=False):
            "Return next n items of the iterable as same type"
            if reset: take.counter = 0
            take.counter += n
            bob = islice(iterable, take.counter-n, take.counter)
            if isinstance(iterable, dict): return dict(bob)
            elif isinstance(iterable, list): return list(bob)
            elif isinstance(iterable, tuple): return tuple(bob)
            elif isinstance(iterable, set): return set(bob)
            elif isinstance(iterable, file): return file(bob)
            else: return bob
        
        def delete_message_batch2(cx, queue, messages): #returns a string reflecting level of success rather than throwing an exception or True/False
          """
          Deletes a list of messages from a queue in a single request.
          :param cx: A boto connection object.
          :param queue: The :class:`boto.sqs.queue.Queue` from which the messages will be deleted
          :param messages: List of any object or structure with id and receipt_handle attributes such as :class:`boto.sqs.message.Message` objects.
          """
          listof10s = []
          asSuc, asErr, acS, acE = "","",0,0
          res = []
          it = tuple(enumerate(messages))
          params = {}
          tenmsg = take(10,it,True)
          while len(tenmsg)>0:
            listof10s.append(tenmsg)
            tenmsg = take(10,it)
          while len(listof10s)>0:
            tenmsg = listof10s.pop()
            params.clear()
            for i, msg in tenmsg: #enumerate(tenmsg):
              prefix = 'DeleteMessageBatchRequestEntry'
              numb = (i%10)+1
              p_name = '%s.%i.Id' % (prefix, numb)
              params[p_name] = msg.get('id')
              p_name = '%s.%i.ReceiptHandle' % (prefix, numb)
              params[p_name] = msg.get('receipt_handle')
            try:
              go = cx.get_object('DeleteMessageBatch', params, BatchResults, queue.id, verb='POST')
              (sSuc,cS),(sErr,cE) = tup_result_messages(go)
              if cS:
                asSuc += ","+sSuc
                acS += cS
              if cE:
                asErr += ","+sErr
                acE += cE
            except cx.ResponseError:
              eprint("Error in batch delete for queue {}({})\nParams ({}) list: {} ".format(queue.name, queue.id, len(params), params))
            except:
              eprint("Error of unknown type in batch delete for queue {}({})\nParams ({}) list: {} ".format(queue.name, queue.id, len(params), params))
          return stringify_final_tup(asSuc, asErr, acS, acE, expect=len(messages)) #mdel #res
        
        def stringify_final_tup(sSuc="", sErr="", cS=0, cE=0, expect=0):
          if sSuc == "": sSuc="None"
          if sErr == "": sErr="None"
          if cS == expect: sSuc="All"
          if cE == expect: sErr="All"
          return "Up to {} messages removed [{}]\t\tMessages remaining ({}) [{}]".format(cS,sSuc,cE,sErr)
        

        【讨论】:

          【解决方案6】:

          类似下面的代码应该可以解决问题。抱歉,它在 C# 中,但转换为 python 应该不难。字典是用来剔除重复的。

              public Dictionary<string, Message> GetAllMessages(int pollSeconds)
              {
                  var msgs = new Dictionary<string, Message>();
                  var end = DateTime.Now.AddSeconds(pollSeconds);
          
                  while (DateTime.Now <= end)
                  {
                      var request = new ReceiveMessageRequest(Url);
                      request.MaxNumberOfMessages = 10;
          
                      var response = GetClient().ReceiveMessage(request);
          
                      foreach (var msg in response.Messages)
                      {
                          if (!msgs.ContainsKey(msg.MessageId))
                          {
                              msgs.Add(msg.MessageId, msg);
                          }
                      }
                  }
          
                  return msgs;
              }
          

          【讨论】:

            猜你喜欢
            • 2012-08-07
            • 2014-10-10
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2013-02-03
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多