【问题标题】:AWS SQS to delay message by X seconds, then next message by X seconds againAWS SQS 将消息延迟 X 秒,然后再将下一条消息延迟 X 秒
【发布时间】:2020-11-13 05:04:43
【问题描述】:

正在寻找一种方法将发送到 lambda 的消息延迟 5 秒。

因此,lambda 收到消息 1,然后 5 秒后收到消息 2,然后 5 秒后收到消息 3,等等,比如说一千条消息。

正在查看 SQS 延迟队列和消息计时器,但它们并不是我想要的。

使用等待的 Step Functions,但在我需要的规模上会很昂贵。

理想情况下,需要一个 SQS 队列,将消息限制为仅每 5 秒发送一次,有没有办法做到这一点?

附言不介意它是 SQS,只需要一个解决方案

【问题讨论】:

  • 您能解释一下为什么 SQS 延迟队列或消息级延迟“不是 [您] 正在寻找的”吗?
  • 是的。使用延迟队列,我可以最终得到,比如说 10 条消息在它们可见之前都被延迟了 5 秒,但是它们都会同时变得可见,所以所有的消息都可以同时被消费,而不是相隔 5 秒。我需要更仔细地查看消息计时器,但从简短的角度来看,它似乎不适用于我的用例
  • 您可以在 lambda 中将 lambda 并发设置为 1 并将批处理大小设置为 1。还要确保 lambda 执行 4-5 秒。这样,消息将在大约 5 秒的时间间隔内从队列中逐一读取。编辑:再想一想,这不是最有效的。
  • 好点,绝对是一个选择。理想情况下,如果我的实际执行时间仅为 1 秒,我不想为 lambda 开启 5 秒付费。
  • 理想情况下,消息的来源应在 5 秒间隔内提交它们。也许带有计划事件的 CloudWatch 日志也会很有用。

标签: amazon-web-services aws-lambda amazon-sqs amazon-sns


【解决方案1】:

您可以使用FIFO Queues。它支持消息排序和每个队列的延迟。了解以下注意事项很重要

  • 消息发送和接收的顺序被严格保留,消息被传递一次并保持可用,直到消费者处理并删除它。
  • FIFO 队列支持允许在单个队列中包含多个有序消息组的消息组。
  • FIFO 队列不支持按消息延迟,仅支持按队列延迟。如果您的应用程序为每条消息设置相同的 DelaySeconds 参数值,则您必须修改您的应用程序以移除每个消息的延迟并改为在整个队列上设置 DelaySeconds。
  • FIFO 队列支持 300 TPS,每个 API 方法(SendMessage、ReceiveMessage 或 DeleteMessage)。您可以使用批处理 API 提取 3000 的 TPS。

现在,如果消息之间的延迟是恒定的,并且您不需要以大于 3000 TPS 的速率填充/排空队列,那么 FIFO 队列可以工作。

【讨论】:

    【解决方案2】:

    您可以使用get_queue_attributes() 并检索“ApproximateNumberOfMessagesDelayed”。这将基本上告诉您队列中当前有多少消息,您可以使用它来乘以所需的延迟时间。为此,您必须单独延迟每条消息,而不是整个队列。 (即DelayTime*ApproximateNumberOfMessagesDelayed + DelayTime)

    【讨论】:

      【解决方案3】:

      我有一个有点类似的问题,但在你的情况下,如果你延迟消息进入队列,那么你不必担心延迟消费消息(在你的情况下是 lambda) .

      正如@Ryan 提到的,

      当您在控制台中发送 Delivery delay(假设为 5 秒)时,它只会延迟整个队列,而 而不是 队列中的单个消息。 这里有一个很好的read 了解Delivery delay

      但诀窍不是延迟队列,而是延迟单个消息(又名aws 术语是Message Queue

      这就是我所做的,

      我先批量发送消息(请阅读docs,到目前为止,您只能批量发送10条消息。)然后为每条消息设置延迟,然后将它们作为batch发送。

      设置发送批量消息

      def setting_to_send_batch_messages(inputDict): 
          """This function sets up a dict. into a batch message(up to 10) so that it can be sent at once (i.e. as a batch) 
      
          Args:
              inputDict ([dict]): [dict. that needs to be batched]
      
          Returns:
              [lst]: [list of dicts of messages]
          """
      
          stock_cnter = 1 # Iterating stock counter
          msg_cnter = 1 # Counter to keep track of number of messages
      
          entryVal_dict = {} # dict. to hold values for each message in the batch
          thisMsgAttribute_perStock_dict = {} # dict. to hold Message Attributes per stock
      
          msg_lst = [] # List to hold all dicts (i.e. stock info) per message
      
          # In the batch, per message delay
          delay_this_message = 0
          # NOTEME: By setting it to 0, means the very first message there is no delay (i.e. sent immediately to the queue) a delay (in seconds) is added to subsequent messages 
      
          # looping over dict.
          for key,val in inputDict.items():
      
              # dict. holding to message attributes
              msgAttributes_dict = {
                  'fieldID' + str(stock_cnter): {
                      'StringValue': key,
                      'DataType': 'String'
                  },
                  'ticker' + str(stock_cnter): {
                      'StringValue': val,
                      'DataType': 'String'
                  }
              }
      
              # By doing an updating, adding to dict. 
              thisMsgAttribute_perStock_dict.update(msgAttributes_dict)
      
              # NOTEME: Per aws sqs, max bumber of MessageAttributes per message is 10, making a message can have only 5 stocks 
              if stock_cnter % 5 == 0 or stock_cnter == len(inputDict): # Checking for 5 stocks OR anything left over grouping by 5
      
                  entryVal_dict['Id'] = str(msg_cnter)
                  entryVal_dict['MessageBody'] =  f'This is the message body for message ID no. {msg_cnter}.'
                  entryVal_dict['DelaySeconds'] = delay_this_message
                  entryVal_dict['MessageAttributes'] = thisMsgAttribute_perStock_dict
      
                  # appending list
                  msg_lst.append(entryVal_dict)
      
                  # resetting dict.
                  entryVal_dict = {}
      
                  delay_this_message += 60 # delaying next message by 1 min.
      
                  msg_cnter += 1 # incrementing message counter
      
                  # resetting dict.
                  thisMsgAttribute_perStock_dict = {}
      
              stock_cnter += 1 # Incrementiing stock loop counter
      
          # print (msg_lst)
          return msg_lst
      

      这是我的inputDict

      {'rec1': 'KO', 'rec0': 'HLT', 'rec2': 'HD', 'rec4': 'AFL', 'rec5': 'STOR', 'rec3': 'WMT',...}
      

      相应地向 SQS 发送批量消息

      def send_sqs_batch_message(entries):    
          # NOTEME: See for more info 
          # https://boto3.amazonaws.com/v1/documentation/api/latest/guide/sqs.html#sending-messages
          sqs_client = boto3.client("sqs", region_name='us-east-2')
      
          response = sqs_client.send_message_batch(
              QueueUrl= YOUR_QUEUE_URL_GOES_HERE,
              Entries = entries
          )
      
          # print(response)
          return response
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2021-10-11
        • 2023-03-29
        • 1970-01-01
        • 2017-02-13
        • 1970-01-01
        • 1970-01-01
        • 2017-09-03
        • 2018-10-11
        相关资源
        最近更新 更多