System Design
B6

Queues and logs

A service hands work or events to others that run later, at their own pace. Pick a queue or a log, decide where the ack goes, and plan for the message that never succeeds.

Not startedSaved in this browser only.
  1. 1A queue deletes a message on ack. A log keeps it, and each consumer group keeps its own offset.
  2. 2Delivery is at least once in practice. An idempotent consumer makes it effectively once.
  3. 3Order holds only inside a partition. Give every message of one entity the same key.
  4. 4Every message needs a way out: a delivery limit and a dead-letter queue.
B6
    A

    Queue or log

    what happens after the read
    Queueone worker per message; the ack deletes itProducersendm6m5m4m3m3, m4: delivered, not ackedWorker A: m3Worker B: m4m1, m2: acked and deletedLogmessages stay; each group keeps its own offsetProducerappends here01234567890, 1: past retention, deleted by agegroup email: 7group search: 4
    • A queue spreads work. Any free worker takes the next message.
    • A log spreads events. Each reader group moves through it at its own speed.
    • A Redis stream with a consumer group has both: per-message acks, and entries that stay.

    A queue hands each message to one worker and forgets it on ack. A log keeps every message, so many groups can read it and replay it.

    B

    Which broker, where

    for background work and events
    methodwherestatus
    Poll a table without SKIP LOCKEDPostgresNot approved
    Table, FOR UPDATE SKIP LOCKEDPostgresJobs with your data
    Pub/Sub channel for jobsRedisNot approved
    List with BRPOPRedisOnly with LMOVE
    Stream with a consumer groupRedisApproved
    Managed queue, standardSQSOrder not needed
    Managed queue, FIFOSQSApproved
    Quorum queue with a delivery limitRabbitMQApproved
    Partitioned log, consumer groupsKafkaReplay, many readers

    Pub/Sub stores nothing: a subscriber that is offline misses the message. BRPOP removes the job before the work, so a crash loses it.

    For jobs that belong to my data, I start with a Postgres table. For events that many services read, I use a log.

    C

    Queue, log, pub/sub

    what each one guarantees
    questionqueuelogpub/sub
    After the ackDeleted.Kept until retention ends. The group's offset moves.Never stored.
    Replay old messagesNot approvedReset the offsetNot approved
    Many readers get every messageOne queue eachOne group eachOnline readers only
    Retry one messagePer messageOffset per partition A stuck message holds back its partition.Not approved
    OrderFIFO queues onlyPer partitionNot approved
    Unit of parallel workA message. Add workers freely.A partition. Members beyond the partition count sit idle.A subscriber.
    ExamplesSQS, RabbitMQ queues, a Postgres tableKafka, Kinesis, Redis streams, RabbitMQ streamsRedis Pub/Sub, WebSocket fan-out

    In a log, one partition is read in order. To retry one message without blocking the rest, copy it to a retry topic and move on.

    I pick by three questions: must the message survive the reader, do several readers need all of it, and must anyone replay it.

    D

    Where the ack goes

    figure, then pseudo code
    At most once: ack, then workreadack↯ crash✕ 0 receipts: lostAt least once: work, then ackreadsend↯ crashread againsendack△ 2 receiptsEffectively once: work and record the id togetherreadid + send↯ crashread againseen: skipack✓ 1 receipt
    The ack is XACK, a row delete, or an offset commit. The crash falls at the same point in each row.
    guaranteehowuse for
    At most onceAck, then work. Never retried.Metrics where a gap is acceptable
    At least onceWork, then ack. Retried until acked.The default for every broker
    Effectively onceAt least once, plus an idempotent handler.Payments, emails, stock
    a worker on a Redis streampseudo code
    worker(name):                        // many copies share one group
      LOOP:
        msgs = XREADGROUP as name, up to 10 new1
        FOR EACH m IN msgs:
          handle(m)
          XACK m                           // only after the effect is durable2
        every few seconds: reclaim()
    
    handle(m):                             // runs as ONE atomic step3
      IF m.job is in the done set: RETURN  // a redelivery: skip it
      add m.job to the done set; send the receipt
    
    reclaim():                             // take over a dead worker's messages
      FOR EACH pending m idle longer than the limit4:
        IF m was delivered 3 times5:
          copy m to the dead letters; XACK m
        ELSE: XAUTOCLAIM m                 // delivery count + 16
    1. 1Each message goes to one member of the group and enters its pending entries list, with a delivery count of 1.
    2. 2Ack last. A crash before this line means a redelivery, never a loss.
    3. 3A server-side script: the check, the mark and the effect cannot be split by a crash.
    4. 4XPENDING lists each pending message with its owner, idle time and delivery count.
    5. 5A poison message fails every time. Stop after 3 tries so it cannot loop forever.
    6. 6XAUTOCLAIM moves the message to this worker and counts one more delivery.
    Tested source Go: read, ack, reclaim, dead letters · Redis script: apply once
    Go: read, ack, reclaim, dead lettersgo
    // Read hands up to n new messages to consumer. Each one enters the pending entries list under
    // that consumer, with a delivery count of 1, and stays there until it is acknowledged.
    func (s *Stream) Read(ctx context.Context, consumer string, n int64) ([]Msg, error) {
      res, err := s.RDB.XReadGroup(ctx, &redis.XReadGroupArgs{
        Group: s.Group, Consumer: consumer, Streams: []string{s.Key, ">"}, Count: n, Block: -1,
      }).Result()
      if errors.Is(err, redis.Nil) {
        return nil, nil
      }
      if err != nil {
        return nil, fmt.Errorf("read for %s: %w", consumer, err)
      }
      var out []Msg
      for _, st := range res {
        for _, m := range st.Messages {
          out = append(out, Msg{ID: m.ID, Job: fmt.Sprint(m.Values["job"]), Deliveries: 1})
        }
      }
      return out, nil
    }
    
    // Ack removes a message from the pending entries list. The entry stays in the stream.
    func (s *Stream) Ack(ctx context.Context, id string) error {
      if err := s.RDB.XAck(ctx, s.Key, s.Group, id).Err(); err != nil {
        return fmt.Errorf("ack %s: %w", id, err)
      }
      return nil
    }
    
    // Reclaim takes over messages that have waited unacknowledged for at least minIdle, for
    // example because their consumer crashed. A message already delivered MaxDeliveries times is
    // a poison message: it moves to the dead-letter stream instead and is acknowledged here.
    func (s *Stream) Reclaim(ctx context.Context, consumer string, minIdle time.Duration, n int64) (claimed, dead []Msg, err error) {
      stale, err := s.RDB.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream: s.Key, Group: s.Group, Idle: minIdle, Start: "-", End: "+", Count: n,
      }).Result()
      if err != nil {
        return nil, nil, fmt.Errorf("list idle pending: %w", err)
      }
      counts := map[string]int64{}
      for _, p := range stale {
        if p.RetryCount < s.MaxDeliveries {
          counts[p.ID] = p.RetryCount
          continue
        }
        m, err := s.deadLetter(ctx, p)
        if err != nil {
          return nil, nil, err
        }
        dead = append(dead, m)
      }
      if len(counts) == 0 {
        return nil, dead, nil
      }
      msgs, _, err := s.RDB.XAutoClaim(ctx, &redis.XAutoClaimArgs{
        Stream: s.Key, Group: s.Group, Consumer: consumer, MinIdle: minIdle, Start: "0-0", Count: n,
      }).Result()
      if err != nil {
        return nil, nil, fmt.Errorf("claim idle pending: %w", err)
      }
      for _, m := range msgs {
        before, ok := counts[m.ID]
        if !ok {
          return nil, nil, fmt.Errorf("claimed %s, which was not idle a moment ago", m.ID)
        }
        claimed = append(claimed, Msg{ID: m.ID, Job: fmt.Sprint(m.Values["job"]), Deliveries: before + 1})
      }
      return claimed, dead, nil
    }
    
    // deadLetter copies a poison message to the dead-letter stream and acknowledges it, in one
    // MULTI/EXEC block, so the message is never in both places or in neither.
    func (s *Stream) deadLetter(ctx context.Context, p redis.XPendingExt) (Msg, error) {
      body, err := s.RDB.XRange(ctx, s.Key, p.ID, p.ID).Result()
      if err != nil || len(body) != 1 {
        return Msg{}, fmt.Errorf("read poison message %s: %w (found %d)", p.ID, err, len(body))
      }
      job := fmt.Sprint(body[0].Values["job"])
      if _, err := s.RDB.TxPipelined(ctx, func(tx redis.Pipeliner) error {
        tx.XAdd(ctx, &redis.XAddArgs{Stream: s.DLQ, Values: []any{"job", job, "source_id", p.ID, "deliveries", p.RetryCount}})
        tx.XAck(ctx, s.Key, s.Group, p.ID)
        return nil
      }); err != nil {
        return Msg{}, fmt.Errorf("move %s to %s: %w", p.ID, s.DLQ, err)
      }
      return Msg{ID: p.ID, Job: job, Deliveries: p.RetryCount}, nil
    }
    
    Redis script: apply oncelua
    -- KEYS[1]: set of job ids already applied. KEYS[2]: hash of receipts sent per job.
    -- ARGV[1]: the job id. Redis runs the whole script as one step.
    if redis.call('SADD', KEYS[1], ARGV[1]) == 0 then
      return 0 -- seen before: a redelivery. Skip the effect.
    end
    redis.call('HINCRBY', KEYS[2], ARGV[1], 1)
    return 1

    I ack after the effect commits, which gives at least once. The handler records the message id with the effect, so a redelivery changes nothing.

    E

    See it: crashes, redelivery and dead letters

    recorded on Redis and Postgres
    Redis stream
    Postgres table
    ProducerXADDRedis stream: jobs, group receiptsorder 1readynot deliveredreceipts: 0order 2readynot deliveredreceipts: 0order 3readynot deliveredreceipts: 0order 4readynot deliveredreceipts: 0entries: 4 · pending: 0Worker ArunningWorker BrunningDead lettersstream jobs:dead0 jobsafter 3 deliveries
    step 1 of 10

    After this step

    entries
    4
    pending
    0
    dead
    0
    receipts
    #1: 0, #2: 0, #3: 0, #4: 0

    Ack first lost a job. Ack last sent one receipt twice. Ack last with a dedupe sent each receipt once, and the bad job went to the dead letters after 3 tries.

    F

    The message path

    click a step; its path lights up
    Order serviceproducerLog or streampartitions, retentionBilling workersgroup billing, NSearch indexergroup searchPostgreseffects, inboxDead lettersafter 3 deliveriesSearch indexits own copy

    Step 1: Produce

    • Write the event with a key, such as the order id.
    • The key picks the partition, so all events of one order stay in order.
    • Retry on a timeout. The broker or the consumer must drop the duplicate.

    If it fails

    The broker is down: retry with backoff. Write the event to an outbox in the same transaction as the data, so a crash loses nothing.

    The producer keys the event, the log keeps it, each group reads at its own offset, and the handler dedupes before it acks.

    G

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    RedisStreams: XADD, XRANGEAn append-only log in memory. Entries stay after an ack until trimmed.Activity feeds, event history
    RedisConsumer groups: XREADGROUP, XACKEach entry goes to one member. The group remembers the last entry it handed out.Job queues
    RedisPending entries list: XPENDINGEvery unacked entry, with its owner, idle time and delivery count.Alerts on stuck work
    RedisXAUTOCLAIM, XCLAIMTake over entries that sat idle. Each claim adds 1 to the delivery count.Worker failover
    RedisScripts and MULTI/EXECDedupe and effect in one step; copy to the dead letters and ack in one step.Any check-then-write
    RedisMemory and asynchronous replicationLimit A failover can lose recent entries. Keep the source of truth in a database.
    PostgresFOR UPDATE SKIP LOCKEDWorkers claim different rows at once. Nobody waits on another worker's lock.Outbox relays, schedulers
    PostgresTransactionsInsert the job in the same transaction as the data it is for.Every multi-row change
    PostgresConditional DELETE and its row countA worker whose claim timed out deletes 0 rows.Fencing, leases
    PostgresUnique key, ON CONFLICT DO NOTHINGThe inbox: a message id takes effect once.Idempotency keys
    KafkaPartitioned, replicated log; retention 7 days by defaultOrder per partition. Many groups replay the same data.Change data capture, event sourcing
    KafkaConsumer groups, committed offsets, rebalancingPartitions move to live members. A new owner resumes at the committed offset.Stream processing
    KafkaCompaction; idempotent producer; transactionsKeep the last value per key. Drop producer retries. Read committed data only.Changelogs, state stores
    SQSVisibility timeout, redrive with maxReceiveCountThe managed form of claim, timeout and dead letters.Background jobs
    RabbitMQPrefetch (basic.qos), quorum queue delivery limitCap unacked messages per consumer. Drop or dead-letter a message after 20 deliveries by default.Routing by topic

    Redis streams give me consumer groups, a pending list with delivery counts, and takeover of idle messages. Postgres gives me SKIP LOCKED and the same transaction as my data.

    H

    A queue in Postgres

    schema, then pseudo code
    jobs and dead jobssql
    -- A job is a row. A claim moves visible_at forward: the visibility timeout.
    -- An ack deletes the row. If the worker dies, the row becomes visible again.
    CREATE TABLE jobs (
      id         bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
      payload    text        NOT NULL,
      attempts   int1         NOT NULL DEFAULT 0,
      visible_at timestamptz NOT NULL DEFAULT now()2,
      locked_by  text3
    );
    CREATE INDEX jobs_visible ON jobs (visible_at, id)4;
    
    -- Jobs that failed max_attempts times wait here for a person or a repair job.
    CREATE TABLE dead_jobs (
      id       bigint PRIMARY KEY,
      payload  text   NOT NULL,
      attempts int    NOT NULL,
      died_at  timestamptz NOT NULL
    );
    1. 1The delivery count. Past the limit, the job moves to dead_jobs.
    2. 2The visibility timeout is a time column. A claim moves it forward; a crash lets it pass.
    3. 3The owner of the current claim. The ack checks it.
    4. 4Workers find visible jobs without scanning the ones in flight.
    claim and ackpseudo code
    claim(worker, n, now):                 // one transaction
      rows = up to n jobs WHERE visible_at <= now,
             oldest first, FOR UPDATE SKIP LOCKED1
      SET visible_at = now + 30 s2, attempts + 1, locked_by = worker
      move rows past 3 attempts3 to dead_jobs
      RETURN the other rows
    
    ack(worker, ids):
      DELETE ids WHERE locked_by = worker4    // a late worker deletes 0 rows
    1. 1Rows another worker is claiming right now are skipped, not waited on.
    2. 2The visibility timeout. Set it longer than the slowest job, or extend it while the job runs.
    3. 3The poison job leaves the queue in the same transaction.
    4. 4A worker whose claim timed out deletes nothing. The job now belongs to someone else.
    Tested source Go: claim, bury, ack · SQL: claim, bury, ack
    Go: claim, bury, ackgo
    // Claim takes up to n visible jobs for worker. A claimed job is hidden until now + Visibility.
    // If the worker does not ack by then, the job is visible again. A job that would start its
    // attempt MaxAttempts + 1 is poison: it moves to dead_jobs in the same transaction.
    func (p *PGQueue) Claim(ctx context.Context, now time.Time, worker string, n int) (jobs, dead []Job, err error) {
      err = pgx.BeginFunc(ctx, p.DB, func(tx pgx.Tx) error {
        rows, err := tx.Query(ctx, q["claim"], now, n, p.Visibility.Milliseconds(), worker)
        if err != nil {
          return fmt.Errorf("claim for %s: %w", worker, err)
        }
        claimed, err := scanJobs(rows)
        if err != nil {
          return fmt.Errorf("claim for %s: %w", worker, err)
        }
        var poison []int64
        for _, j := range claimed {
          if j.Attempts > p.MaxAttempts {
            poison = append(poison, j.ID)
            dead = append(dead, Job{ID: j.ID, Payload: j.Payload, Attempts: j.Attempts - 1})
            continue
          }
          jobs = append(jobs, j)
        }
        if len(poison) > 0 {
          if _, err := tx.Exec(ctx, q["bury"], poison, now); err != nil {
            return fmt.Errorf("bury %v: %w", poison, err)
          }
        }
        return nil
      })
      return jobs, dead, err
    }
    
    // Ack deletes finished jobs that worker still owns and returns how many it deleted.
    func (p *PGQueue) Ack(ctx context.Context, worker string, ids ...int64) (int64, error) {
      tag, err := p.DB.Exec(ctx, q["ack"], ids, worker)
      if err != nil {
        return 0, fmt.Errorf("ack %v for %s: %w", ids, worker, err)
      }
      return tag.RowsAffected(), nil
    }
    
    SQL: claim, bury, acksql
    -- Take up to $2 visible jobs. SKIP LOCKED passes over rows another worker is claiming right now,
    -- and the new visible_at hides the claimed rows from every worker until the timeout ends.
    UPDATE jobs
    SET visible_at = $1 + $3 * interval '1 millisecond',
        attempts   = attempts + 1,
        locked_by  = $4
    WHERE id IN (
      SELECT id FROM jobs
      WHERE visible_at <= $1
      ORDER BY id
      LIMIT $2
      FOR UPDATE SKIP LOCKED
    )
    RETURNING id, payload, attempts;
    
    -- A claimed job past max_attempts is poison. Move it aside instead of handing it out again.
    WITH dead AS (
      DELETE FROM jobs WHERE id = ANY($1)
      RETURNING id, payload, attempts - 1 AS deliveries
    )
    INSERT INTO dead_jobs (id, payload, attempts, died_at)
    SELECT id, payload, deliveries, $2 FROM dead;
    
    -- Delete only while this worker still owns the claim. After a timeout, another worker owns it.
    DELETE FROM jobs WHERE id = ANY($1) AND locked_by = $2;
    • Delete on ack keeps the table small. VACUUM removes the dead rows.
    • Many idle workers that poll add load. Use LISTEN and NOTIFY, or a short sleep when empty.

    A worker claims rows with SKIP LOCKED and hides them with a visibility timeout. The ack deletes the row only while the worker still owns it.

    I

    Partitions and consumer groups

    the unit of order and of parallel work
    key · sequencepartition = hash(key) mod 4Producerkey = user idu7 → P2P0a·1c·1a·2c·2P1b·1d·1b·2P2u7·1e·1u7·2u7·3P3f·1g·1f·2group billinggroup searchc1P0, P1c2P2, P3s1P0 to P3u7·1, u7·2, u7·3 sit in P2, in order, and only c2 reads them in group billing.
    • More partitions allow more parallel members, but adding partitions moves keys to new partitions.
    • One hot key stays in one partition. Its owner is the limit for that key.

    The key picks the partition, so one user's events stay in order. Inside a group, one partition has one owner.

    J

    Rebalancing and lag

    what moves partitions, and what it costs
    eventwhat happens
    A member joinsPartitions move to it. In an eager rebalance, every member pauses.
    A member leaves cleanlyIt commits its offsets first. No message runs twice.
    A member crashesIts partitions stall until the session timeout. The new owner replays from the last commit.
    A member stops pollingPast max.poll.interval.ms, 5 minutes by default, the group removes it.
    More members than partitionsThe extra members own nothing.
    • Lag is the newest offset minus the committed offset, per partition.
    • Lag divided by the handle rate gives lag in seconds. Alert on that.
    • Cooperative rebalancing moves only the partitions that change owner, so the others keep working.

    A crash costs a session timeout of stalled partitions plus a replay since the last commit. I alert on lag measured in time.

    K

    See it: a partitioned log

    seeded simulation, 6 partitions, 1 tick per second
    partition by
    jump to
    lag: messages in the log not yet handledproduce 40/sproduce 100/sproduce 50/s02505007501000020406080100120 sP0c1P1c1c3P2c1c2c3c4P3c2c3c5P4c2c3c4c6P5c2c3c4c7owner of each partitiongrey: rebalance pause · red strip: owner crashed · marks under the axis: duplicates (amber), out of order (red)
    at 76 s
    produce 50/s; 3 members, none idle; paused for a rebalance
    lag
    330 messages
    this second
    0 handled, 0 replayed, 0 out of order
    whole run
    7250 produced; 60 replayed after the crash
    peak lag
    880 messages
    order
    ✓ every user’s messages handled in order

    Each member handles 30 messages a second and commits every 5 s. The group notices a crash after 4 s; a rebalance pauses the group for 2 s. The crash replayed 60 messages.

    With a key, every user's messages stayed in order, even through a crash. Without a key, 684 messages ran after newer ones of the same user.

    L

    Retention and compaction

    how long a log keeps data
    before: every update, in orderk1=aoffset 0k2=xoffset 1k1=boffset 2k3=poffset 3k2=yoffset 4k1=coffset 5k3=nulloffset 6k2=zoffset 7after compaction: the last value of each keyk1=coffset 5k3=nulloffset 6k2=zoffset 7k3=null is a tombstone: kept 1 day (delete.retention.ms), then k3 is gone.A new consumer replays it to rebuild the latest value of every key.
    policykeepsuse for
    Delete by time or sizeThe last N days. Kafka: 7 days by default.Events, replay after a bug
    CompactThe last value of each keyChangelogs, current state
    Trim by length or idThe newest entries (XADD MAXLEN or MINID)Redis streams in memory
    Delete on ackUnacked messages only. SQS: 4 days by default.Work queues

    Delete retention keeps a time window for replay. Compaction keeps the last value per key, so the topic works as a table.

    M

    Backpressure

    when consumers are slower than producers
    methodeffectstatus
    Pull, with a batch limitEach consumer takes only what it can handle: COUNT, max.poll.records (500 by default), prefetch.Approved
    Scale consumers on lagAdd members up to the partition count.Approved
    Bounded buffer in the producerWhen full, block or return 429. The caller slows down.Approved
    Trim by lengthCaps memory, but deletes unread entries.Loss is acceptable
    Unbounded in-memory queueGrows until the process runs out of memory.Not approved
    • A queue that is always growing means the consumers are too slow for the rate. A buffer only delays the failure.
    • Retry with backoff and jitter, or a failing downstream gets retries from every consumer at once.

    A log buffers on disk, so the producer keeps going while lag grows. I bound the work in flight per consumer and scale consumers on lag.

    N

    Failure cases

    what breaks, and what keeps the work right
    eventresultwhy it is safesaved by
    A worker crashes before the ackThe message stays pending, or hidden.After the idle limit or the timeout, another worker takes it. The inbox absorbs the repeat.Pending list
    A worker acks, then crashesThe job is lost.Not safe. Ack only after the effect commits.Ack last
    A slow worker passes the visibility timeoutTwo workers run the job.The late ack deletes 0 rows. The handler is idempotent.locked_by
    A message fails every timeIt comes back after each timeout.After 3 deliveries it moves to the dead letters. A person or a repair job looks at it.Delivery count
    A group member crashesIts partitions stall, then replay since the last commit.The new owner starts at the committed offset; the inbox drops the replays.Offsets
    Members join and leave again and againEach rebalance pauses the group.Use cooperative rebalancing. Do not autoscale on every spike.Cooperative
    One key is far hotter than the restOne partition lags; the others are idle.Split the key, such as user id plus a bucket, where order across buckets does not matter.Key design
    Lag passes the retentionEntries are deleted before the group reads them.Not safe. Alert when lag in time nears the retention.Lag alert
    The producer retries after a timeoutThe same event is appended twice.The idempotent producer drops it; the consumer inbox catches the rest.Inbox
    Redis fails overRecent entries can be lost.The outbox in Postgres still has the event. The relay sends it again.Outbox
    O

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Postgres table2+ Redis stream3managed queue4partitioned logmore load →
    Consume rates and demand1001k10k100k1MDemand: 100M jobs a day: 1,157 messages per secondDemand: 100M jobs a day1,157Demand: 10× peak: 11,574 messages per secondDemand: 10× peak11,574Postgres, batch 1: 2,973 messages per secondPostgres, batch 12,973Postgres, batch 100: 82,970 messages per secondPostgres, batch 10082,970Redis stream, batch 1: 11,478 messages per secondRedis stream, batch 111,478Redis stream, batch 100: 231,826 messages per secondRedis stream, batch 100231,826SQS FIFO, batches of 10: 3k per API actionSQS FIFO, batches of 103k per API actionKinesis, one shard: 1k records writtenKinesis, one shard1k records writtenmessages per second, log scale
    stepaddit handlesmove up when you see
    1A Postgres table with SKIP LOCKED, a visibility timeout and a dead_jobs table.About 3,000 jobs a second claimed one at a time, 83,000 in batches of 100. Jobs commit with the data.Queue traffic competes with the main workload, or several services need the same events.
    2A Redis stream with consumer groups, fed by an outbox.About 11,000 a second one at a time, 232,000 in batches of 100, on one Redis.The stream must survive a node loss, or outgrow memory.
    3A managed queue with a redrive policy.Standard queues: nearly unlimited calls a second. FIFO: 300 calls a second per action, 3,000 messages with batches of 10.Many groups must read and replay the same events, in order per key.
    4A partitioned log, keyed by entity, with replication and retention.Each partition adds parallel work; each group adds a reader. Kinesis: 1,000 records or 1 MB a second written per shard.Top of the ladder. Add partitions before the members run out.

    Demand example: 100 million jobs a day is 1,157 a second; a 10× peak is 11,574. Rates are 8 clients on one shared laptop and varied by up to 4× between runs, so read them as orders of magnitude. SQS and Kinesis figures are from the AWS quotas.

    I start with a Postgres table and add a Redis stream for rate. When many readers need replay in key order, I move to a partitioned log.

    P

    Drill

    predict, then reveal

    0 of 10 known

    1. A worker acks a message, then crashes before it sends the receipt. What does the customer get?

    2. A worker sends the receipt, then crashes before XACK. What happens next?

    3. How do you get effectively once on top of at least once?

    4. You have 6 partitions and start 8 consumers in one group. What happens?

    5. Events of one order must be handled in order. How do you publish them?

    6. One message fails every time it is handled. What stops it from blocking or looping forever?

    7. A member crashes. Why do some messages run twice, when a member that leaves cleanly causes none?

    8. Consumer lag grows for an hour. What do you check?

    9. A slow worker passes its visibility timeout. Another worker takes the job. What protects you?

    10. When is a Postgres table the right queue?

    Q

    Numbers to say

    measured, then documented
    Postgres
    About 3,000 jobs a second at batch 1; 83,000 at batch 100.
    Redis
    About 11,000 messages a second at batch 1; 232,000 at batch 100.
    produce
    Redis 267,000, Postgres 249,000 a second, batches of 100.
    Kafka
    Retention 7 days. Session timeout 45 s. max.poll.records 500.
    SQS
    Visibility timeout 30 s, up to 12 h. Retention 4 days, up to 14. Messages up to 1 MiB.
    RabbitMQ
    Quorum queues drop or dead-letter a message after 20 deliveries by default.

    Postgres 16 and Redis 8 on an 8-core laptop, 8 clients, median of 3 runs of 2 s. Other work shared the machine, and repeat runs varied by up to 4×. Defaults are from the Kafka, AWS and RabbitMQ documentation.