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.
- 1A queue deletes a message on ack. A log keeps it, and each consumer group keeps its own offset.
- 2Delivery is at least once in practice. An idempotent consumer makes it effectively once.
- 3Order holds only inside a partition. Give every message of one entity the same key.
- 4Every message needs a way out: a delivery limit and a dead-letter queue.
- 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.
| method | where | status |
|---|---|---|
| Poll a table without SKIP LOCKED | Postgres | Not approved |
| Table, FOR UPDATE SKIP LOCKED | Postgres | Jobs with your data |
| Pub/Sub channel for jobs | Redis | Not approved |
| List with BRPOP | Redis | Only with LMOVE |
| Stream with a consumer group | Redis | Approved |
| Managed queue, standard | SQS | Order not needed |
| Managed queue, FIFO | SQS | Approved |
| Quorum queue with a delivery limit | RabbitMQ | Approved |
| Partitioned log, consumer groups | Kafka | Replay, 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.
| question | queue | log | pub/sub |
|---|---|---|---|
| After the ack | Deleted. | Kept until retention ends. The group's offset moves. | Never stored. |
| Replay old messages | Not approved | Reset the offset | Not approved |
| Many readers get every message | One queue each | One group each | Online readers only |
| Retry one message | Per message | Offset per partition A stuck message holds back its partition. | Not approved |
| Order | FIFO queues only | Per partition | Not approved |
| Unit of parallel work | A message. Add workers freely. | A partition. Members beyond the partition count sit idle. | A subscriber. |
| Examples | SQS, RabbitMQ queues, a Postgres table | Kafka, Kinesis, Redis streams, RabbitMQ streams | Redis 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.
| guarantee | how | use for |
|---|---|---|
| At most once | Ack, then work. Never retried. | Metrics where a gap is acceptable |
| At least once | Work, then ack. Retried until acked. | The default for every broker |
| Effectively once | At least once, plus an idempotent handler. | Payments, emails, stock |
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- 1Each message goes to one member of the group and enters its pending entries list, with a delivery count of 1.
- 2Ack last. A crash before this line means a redelivery, never a loss.
- 3A server-side script: the check, the mark and the effect cannot be split by a crash.
- 4XPENDING lists each pending message with its owner, idle time and delivery count.
- 5A poison message fails every time. Stop after 3 tries so it cannot loop forever.
- 6XAUTOCLAIM moves the message to this worker and counts one more delivery.
Tested source Go: read, ack, reclaim, dead letters · Redis script: apply once
// 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
}
-- 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 1I ack after the effect commits, which gives at least once. The handler records the message id with the effect, so a redelivery changes nothing.
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.
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.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Redis | Streams: XADD, XRANGE | An append-only log in memory. Entries stay after an ack until trimmed. | Activity feeds, event history |
| Redis | Consumer groups: XREADGROUP, XACK | Each entry goes to one member. The group remembers the last entry it handed out. | Job queues |
| Redis | Pending entries list: XPENDING | Every unacked entry, with its owner, idle time and delivery count. | Alerts on stuck work |
| Redis | XAUTOCLAIM, XCLAIM | Take over entries that sat idle. Each claim adds 1 to the delivery count. | Worker failover |
| Redis | Scripts and MULTI/EXEC | Dedupe and effect in one step; copy to the dead letters and ack in one step. | Any check-then-write |
| Redis | Memory and asynchronous replication | Limit A failover can lose recent entries. Keep the source of truth in a database. | |
| Postgres | FOR UPDATE SKIP LOCKED | Workers claim different rows at once. Nobody waits on another worker's lock. | Outbox relays, schedulers |
| Postgres | Transactions | Insert the job in the same transaction as the data it is for. | Every multi-row change |
| Postgres | Conditional DELETE and its row count | A worker whose claim timed out deletes 0 rows. | Fencing, leases |
| Postgres | Unique key, ON CONFLICT DO NOTHING | The inbox: a message id takes effect once. | Idempotency keys |
| Kafka | Partitioned, replicated log; retention 7 days by default | Order per partition. Many groups replay the same data. | Change data capture, event sourcing |
| Kafka | Consumer groups, committed offsets, rebalancing | Partitions move to live members. A new owner resumes at the committed offset. | Stream processing |
| Kafka | Compaction; idempotent producer; transactions | Keep the last value per key. Drop producer retries. Read committed data only. | Changelogs, state stores |
| SQS | Visibility timeout, redrive with maxReceiveCount | The managed form of claim, timeout and dead letters. | Background jobs |
| RabbitMQ | Prefetch (basic.qos), quorum queue delivery limit | Cap 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.
-- 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
);- 1The delivery count. Past the limit, the job moves to dead_jobs.
- 2The visibility timeout is a time column. A claim moves it forward; a crash lets it pass.
- 3The owner of the current claim. The ack checks it.
- 4Workers find visible jobs without scanning the ones in flight.
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- 1Rows another worker is claiming right now are skipped, not waited on.
- 2The visibility timeout. Set it longer than the slowest job, or extend it while the job runs.
- 3The poison job leaves the queue in the same transaction.
- 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
// 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
}
-- 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.
- 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.
| event | what happens |
|---|---|
| A member joins | Partitions move to it. In an eager rebalance, every member pauses. |
| A member leaves cleanly | It commits its offsets first. No message runs twice. |
| A member crashes | Its partitions stall until the session timeout. The new owner replays from the last commit. |
| A member stops polling | Past max.poll.interval.ms, 5 minutes by default, the group removes it. |
| More members than partitions | The 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.
- 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.
| policy | keeps | use for |
|---|---|---|
| Delete by time or size | The last N days. Kafka: 7 days by default. | Events, replay after a bug |
| Compact | The last value of each key | Changelogs, current state |
| Trim by length or id | The newest entries (XADD MAXLEN or MINID) | Redis streams in memory |
| Delete on ack | Unacked 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.
| method | effect | status |
|---|---|---|
| Pull, with a batch limit | Each consumer takes only what it can handle: COUNT, max.poll.records (500 by default), prefetch. | Approved |
| Scale consumers on lag | Add members up to the partition count. | Approved |
| Bounded buffer in the producer | When full, block or return 429. The caller slows down. | Approved |
| Trim by length | Caps memory, but deletes unread entries. | Loss is acceptable |
| Unbounded in-memory queue | Grows 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.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A worker crashes before the ack | The 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 crashes | The job is lost. | Not safe. Ack only after the effect commits. | Ack last |
| A slow worker passes the visibility timeout | Two workers run the job. | The late ack deletes 0 rows. The handler is idempotent. | locked_by |
| A message fails every time | It 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 crashes | Its 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 again | Each rebalance pauses the group. | Use cooperative rebalancing. Do not autoscale on every spike. | Cooperative |
| One key is far hotter than the rest | One 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 retention | Entries are deleted before the group reads them. | Not safe. Alert when lag in time nears the retention. | Lag alert |
| The producer retries after a timeout | The same event is appended twice. | The idempotent producer drops it; the consumer inbox catches the rest. | Inbox |
| Redis fails over | Recent entries can be lost. | The outbox in Postgres still has the event. The relay sends it again. | Outbox |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | A 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. |
| 2 | A 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. |
| 3 | A 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. |
| 4 | A 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.
0 of 10 known
A worker acks a message, then crashes before it sends the receipt. What does the customer get?
A worker sends the receipt, then crashes before XACK. What happens next?
How do you get effectively once on top of at least once?
You have 6 partitions and start 8 consumers in one group. What happens?
Events of one order must be handled in order. How do you publish them?
One message fails every time it is handled. What stops it from blocking or looping forever?
A member crashes. Why do some messages run twice, when a member that leaves cleanly causes none?
Consumer lag grows for an hour. What do you check?
A slow worker passes its visibility timeout. Another worker takes the job. What protects you?
When is a Postgres table the right queue?
- 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.