Optimizing SQS batching for groups

Viewed 367

I have a system built on AWS which uses SQS to put requests related to different customers in a queue, to be handle by a single consumer - a Lambda function.

One special circumstance is that requests should be handled in such a way that each lambda invocation should handle requests from one customer only. To achieve this separation, the lambda function receives SQS messages, groups them by customer ID, and then invokes itself once per customer (if there were more than 1 customer present in the batch of messages).

Example: Consider a queue of items CxMy (message y from customer x):

[ C1M1, C2M1, C1M2, C3M1, C1M2, C3M2, C2M2, C2M2, C3M3, C1M3, C4M1, C3M2 ]

(Note that messages might be duplicated as well.)

This set of records invoke the lambda function, which groups messages like this:

C1: [M1, M2, M3]
C2: [M1, M2]
C3: [M1, M2, M3]
C4: [M1]

So we end up with the Lambda function invoking itself 4 times.

Now, I also want to utilize buffering, so that Lambda invocations may be delayed up to 30 seconds to fill up more records at a time, so I use ReceiveMessageWaitTimeSeconds: 30 on the queue, and BatchSize: 50 and MaximumBatchingWindowInSeconds: 30 on the lambda function.

My problem is that invocations aren't really optimized, one invocation could have messages from 10 different customers, leading to a lot of extra splitting up of records and recursive invocations. My ideal scenario would be if messages from each customer are grouped together to as large extent as possible before invoking the Lambda function.

I have looked at FIFO queues and the message group feature, but I'm not sure how much this would help.

  1. I don't need deduplication since I do this in the application layer anyway, and I don't want to ignore processing a message 5 seconds after the same message was received, if it has already been processed. In other words, if I receive M2 5 times in a batch I will just process M2 once, but if I get it once in one batch and once in another batch 5 seconds later, I do need to process it both times. Using deduplication, I fear I wouldn't get the second message?
  2. I don't need ordering
  3. FIFOs have a message group concept, but it seems mostly used to guarantee order of delivery within each group, not to group messages from the same group into the same batch? There are some comments about how "if possible", messages from the same group will be put in the same batch, but I'm not sure to what extent this helps me?
  4. FIFO's have lower throughput

My messages are basically events saying some data has changed, meaning I need to read some data after being notified about an event. If I get 10 events about the same data in a short amount of time, ideally I only want to refresh the data once, but I do need to ensure that I always refresh the data after the most recent message was received, so I don't want to throw away any messages - only hold off refreshing for a few seconds to have a chance to buffer up duplicates before refreshing the data.

What are my best options in optimizing this system? I have also thought about actually creating separate SQS queues. The number of customers is finite, let's say <100, but could be created dynamically so I would have to support creating new SQS queues in runtime upon detecting the addition of a new customer. I guess separate queues would be ideal for grouping messages, but it also seems quite complex to create new AWS resources from my backend ad-hoc. Is there a good way I could use features of SQS to achieve what I want without resorting to separate queues?

0 Answers
Related