FIFO Ordering with SQS Message Group IDs

This guide puts the per-key lane model from Message Ordering Guarantees into practice on AWS, as part of Queue Fundamentals & Architecture. It configures an SQS FIFO queue so that every message for one customer account is processed strictly in sequence while thousands of accounts proceed in parallel, and it covers the consumer mistakes that quietly break that guarantee.

Problem Statement

A billing service publishes ChargeCreated, ChargeCaptured, and ChargeRefunded events to a standard SQS queue, and an accounting worker applies each one to a ledger. A handful of times a week the ledger shows a refund applied to a charge that does not exist yet, because a standard queue makes no ordering promise and two workers raced. You need strict order for each account's events, parallelism across accounts, no duplicate application when producers retry, and poison events that do not block an account forever.

Prerequisites

  • AWS CLI v2 or Terraform 1.4+, with sqs:* on the target queues.
  • A queue name ending in .fifo — the suffix is mandatory for FIFO queues and cannot be added later; migrating from a standard queue means creating a new queue.
  • A stable identifier for the ordering scope in every message (here account_id) and a unique id for every event (here event_id).
  • A consumer that can be made to stop mid-batch on failure (all AWS SDKs can; some wrapper libraries hide this).

Step 1 — Create the FIFO Queue and Its DLQ

Create the dead-letter queue first, because the source queue's redrive policy references its ARN. Both must be FIFO.

# sqs_fifo.tf
resource "aws_sqs_queue" "ledger_dlq" {
  name                      = "ledger-events-dlq.fifo"
  fifo_queue                = true
  message_retention_seconds = 1209600          # 14 days to triage a parked event
}

resource "aws_sqs_queue" "ledger" {
  name                        = "ledger-events.fifo"
  fifo_queue                  = true
  content_based_deduplication = false          # we supply MessageDeduplicationId explicitly
  deduplication_scope         = "messageGroup" # required for high-throughput mode
  fifo_throughput_limit       = "perMessageGroupId"
  visibility_timeout_seconds  = 60             # > p99 handler time for a full batch
  receive_wait_time_seconds   = 20             # long polling by default
  redrive_policy = jsonencode({
    deadLetterTargetArn = aws_sqs_queue.ledger_dlq.arn
    maxReceiveCount     = 5                    # after 5 receives the head is parked
  })
}

deduplication_scope = "messageGroup" plus fifo_throughput_limit = "perMessageGroupId" enables high-throughput FIFO, which raises the ceiling from 300 API calls per second per queue to a per-group limit with a much higher queue-wide total. Enable it at creation; there is no reason to start in the legacy mode for a new queue.

One in-flight message per group Three message groups sit in the FIFO queue. The head of each group can be in flight at the same time, so three accounts progress in parallel, but the second message of any group is not delivered until the first is deleted. Groups run in parallel, messages within a group do not acct-17 in flight waiting waiting acct-42 in flight waiting acct-90 in flight SQS enforces this next message in a group is hidden until the head is deleted Throughput scales with the number of groups that have work, not with the number of consumers.

Step 2 — Choose the MessageGroupId

The group id is the ordering scope. It should be the smallest entity whose events must be applied in sequence. For a ledger that is the account: charges on different accounts are independent, but a refund on account 17 depends on the charge on account 17.

# group_ids.py — derive the ordering key in one place so every producer agrees
def ledger_group_id(event: dict) -> str:
    # One lane per account. Never use a constant ("ledger") — that serialises
    # the whole queue through one in-flight message.
    return f"acct-{event['account_id']}"

Two anti-patterns are common. Using a constant group id gives global order and a throughput of one message at a time. Using a random or per-message group id gives no ordering at all and simply pays the FIFO price for nothing. If you find the "right" key is a tenant with millions of events per day, ordering per tenant will make that tenant's lane the bottleneck; look for a finer key (per account within a tenant) that still covers every dependency.

Step 3 — Publish with Explicit Deduplication IDs

FIFO queues discard a message whose deduplication id was already seen in the last five minutes. Set it to the event id so a producer retry after a timeout does not enqueue the event twice.

# producer.py
import json
import boto3

sqs = boto3.client("sqs")
QUEUE_URL = "https://sqs.eu-west-1.amazonaws.com/123456789012/ledger-events.fifo"

def publish(event: dict) -> None:
    sqs.send_message(
        QueueUrl=QUEUE_URL,
        MessageBody=json.dumps(event),
        MessageGroupId=ledger_group_id(event),
        MessageDeduplicationId=event["event_id"],   # stable across producer retries
    )

def publish_batch(events: list[dict]) -> None:
    # Up to 10 per call; order is preserved within each group in the batch.
    entries = [{
        "Id": str(i),
        "MessageBody": json.dumps(e),
        "MessageGroupId": ledger_group_id(e),
        "MessageDeduplicationId": e["event_id"],
    } for i, e in enumerate(events)]
    resp = sqs.send_message_batch(QueueUrl=QUEUE_URL, Entries=entries)
    if resp.get("Failed"):
        # A partial failure can reorder a group if you retry only the failed entry
        # after later entries succeeded. Retry the failed ones BEFORE sending more
        # events for the same group.
        raise RuntimeError(f"batch partially failed: {resp['Failed']}")

The comment on partial batch failure is the subtle part: if entry 3 of a batch fails and entries 4–10 succeed, and entry 4 is for the same group, a naive retry of entry 3 puts it after entry 4. Keep each group's events in a single batch, or stop publishing for a group until its failed entries are resent.

Step 4 — Consume Without Skipping a Failed Head

A single ReceiveMessage call can return up to ten messages, and several can belong to the same group, in order. The consumer must process them sequentially and stop at the first failure for that group.

# consumer.py
import json
import boto3
from collections import defaultdict

sqs = boto3.client("sqs")

def run() -> None:
    while True:
        resp = sqs.receive_message(
            QueueUrl=QUEUE_URL,
            MaxNumberOfMessages=10,
            WaitTimeSeconds=20,
            AttributeNames=["MessageGroupId", "ApproximateReceiveCount"],
        )
        by_group: dict[str, list] = defaultdict(list)
        for msg in resp.get("Messages", []):
            by_group[msg["Attributes"]["MessageGroupId"]].append(msg)

        for group, msgs in by_group.items():
            for msg in msgs:                       # already in group order
                try:
                    apply_to_ledger(json.loads(msg["Body"]))
                except Exception:
                    # Do NOT delete, do NOT continue with this group. The rest of
                    # the group stays invisible until the visibility timeout ends,
                    # then SQS redelivers starting from this message.
                    break
                sqs.delete_message(QueueUrl=QUEUE_URL, ReceiptHandle=msg["ReceiptHandle"])

Processing different groups concurrently (one task per group) is safe and a good way to use a multi-core worker; processing messages within a group concurrently is not.

Stop the group at the first failure A batch of messages is grouped by message group. Group A processes three messages successfully. In group B the first message fails, so the consumer stops and leaves the second message undeleted; both are redelivered together, in order, after the visibility timeout. One receive batch, two groups group A A1 applied, deleted A2 applied, deleted A3 applied, deleted group B B1 raises B2 skipped, not deleted both return after the timeout Applying B2 here would record a capture for a charge the ledger has not seen.

Step 5 — Decide What Happens When the Head Keeps Failing

With maxReceiveCount = 5, a head message that fails five times moves to ledger-events-dlq.fifo, and the group continues with the next message. For a ledger, continuing past a missing charge is wrong, so the consumer must detect the gap. The cleanest approach is a per-account sequence number assigned by the producer, checked by the consumer before applying:

# gap_check.py
def apply_to_ledger(event: dict) -> None:
    last = ledger.last_sequence(event["account_id"])
    if event["seq"] <= last:
        return                                   # duplicate delivery: already applied
    if event["seq"] != last + 1:
        # A predecessor was parked in the DLQ. Refuse to apply; this message will
        # follow it to the DLQ after maxReceiveCount, keeping the account frozen
        # rather than wrong. Alert so someone replays the missing event first.
        raise GapDetected(event["account_id"], expected=last + 1, got=event["seq"])
    ledger.apply(event)                          # writes entry and seq atomically

This turns the automatic "skip and park" behaviour into "block and park": every later event for the broken account follows the poison message into the DLQ, and a replay in order restores the account. The pattern is covered in depth in handling out-of-order events with sequence numbers.

Step 6 — Size the Visibility Timeout for a Whole Group Batch

Because a consumer may process up to ten messages of one group sequentially, the visibility timeout must cover the batch, not one message. If the timeout expires mid-group, SQS makes the remaining messages visible again and another consumer may start on message 6 while the first is still on message 5.

# Extend visibility while working through a long group batch
def extend_if_needed(msgs: list, started: float, budget: int = 60) -> None:
    if time.monotonic() - started > budget * 0.6:
        sqs.change_message_visibility_batch(
            QueueUrl=QUEUE_URL,
            Entries=[{"Id": str(i), "ReceiptHandle": m["ReceiptHandle"],
                      "VisibilityTimeout": budget} for i, m in enumerate(msgs)],
        )

The general rules are in the visibility timeout deep dive; for FIFO the unit of work is the group batch.

Verification

Publish a sequence for one account with deliberate handler delays and confirm the ledger applies them in order:

# Send 50 sequenced events for one account and 50 for another
python -c '
from producer import publish
for i in range(1, 51):
    publish({"event_id": f"a17-{i}", "account_id": 17, "seq": i, "type": "Adjust", "amount": 1})
    publish({"event_id": f"a42-{i}", "account_id": 42, "seq": i, "type": "Adjust", "amount": 1})
'
# Afterwards, every account's applied sequence must be 1..50 with no gaps or repeats
psql -c "select account_id, array_agg(seq order by applied_at) = array(select generate_series(1,50))
         from ledger_entries where account_id in (17,42) group by account_id;"

Both rows should return t. Then check the throughput side: with 200 distinct accounts in flight and four consumer processes, NumberOfMessagesDeleted in CloudWatch should rise roughly with the number of consumers, while a test with one account should stay flat no matter how many consumers run.

Gotchas & Edge Cases

The five-minute deduplication window. Deduplication only covers five minutes. A producer that retries after a longer outage, or a replay from the DLQ, can apply the same event twice — the sequence-number check in Step 5 is what makes the consumer idempotent beyond that window, in line with idempotent job design.

Lambda triggers and batch failures. When SQS FIFO triggers AWS Lambda, a function error returns the whole batch. With ReportBatchItemFailures, you must report the failed message and every later message of the same group as failed, or Lambda deletes them and breaks order.

In-flight limit. A FIFO queue allows at most 20,000 in-flight messages. A consumer fleet that receives far more than it can process in one visibility window hits the limit and receives empty responses while the queue looks full.

Hot groups. One very active group runs at the speed of one sequential consumer. ApproximateAgeOfOldestMessage will look healthy on average while one account lags by minutes; emit a per-group lag metric from the consumer.

Group id granularity A constant group id orders everything and allows one message in flight. A per-account id orders each account and runs accounts in parallel. A random id per message allows full parallelism but provides no ordering. Too coarse, right, too fine "ledger" (constant) global order 1 message in flight "acct-17" (per entity) order where it matters accounts in parallel uuid4() per message no ordering at all FIFO cost, no benefit

FAQ

Can I convert an existing standard SQS queue to FIFO? No. The queue type is fixed at creation and FIFO queue names must end in .fifo. Create the new FIFO queue, point producers at it, let consumers drain the old standard queue, and then remove it. During the cutover, events for the same entity can exist in both queues, so either pause producers briefly or make the consumer's sequence check tolerate the overlap.

How many message groups can a FIFO queue have? There is no practical limit on the number of groups; millions of distinct group ids are normal. What is limited is the number of in-flight messages (20,000 per queue) and, in high-throughput mode, the request rate per group. Many small groups is the healthy shape.

Does MessageGroupId affect which consumer gets the message? Not in a sticky way. SQS does not pin a group to one consumer; any consumer can receive the next message of a group once the previous one is deleted. That is why the consumer must not keep per-group state in memory between receives — load it from the database each time.

Should I use content-based deduplication instead of explicit IDs? Only if identical bodies always mean the same event. Content-based deduplication hashes the body, so two genuinely different events with the same payload (two identical $5 adjustments) would be collapsed into one. An explicit event id is safer for anything financial.

Related