You Might Not Need Redis or Kafka Yet, Postgres as a Cache and a Message Queue
A slow endpoint and an email that should not block a request are two problems almost every backend meets. The reflex is to add Redis and a message broker. This is how to solve both with the Postgres you already run, how it works underneath, and an honest comparison with Redis, Memcached, Kafka, RabbitMQ, SQS and NATS.
- postgres
- caching
- queue
- architecture
- redis
- kafka
- golang
## Two problems every backend meets
Most backends start the same way: an application and a database. Then, sooner or later, two things happen.
The first is a slow endpoint. A page that lists courses, products or reports runs a query with several joins and a count. It takes a few hundred milliseconds, it is requested constantly, and its answer barely changes from one minute to the next. Running the same expensive query again and again is plainly wasteful.
The second is work that should not happen inside a request. A user registers and you send a confirmation email. The email provider is slow today, so the registration request hangs. Or the provider is down for ten seconds, the send fails, and the user is registered but never told.
Both problems have well-known answers. For the first you add a cache, and the default choice is Redis. For the second you add a queue, and the default choice is a message broker such as RabbitMQ or Kafka.
Those answers are correct. They also turn one moving part into three. Each new service has to be deployed, secured, monitored, upgraded, backed up, and run on every developer's laptop. For a large system that cost is small next to the benefit. For a small or medium one, it is often the biggest thing on the bill.
On a platform I worked on recently, the application shipped as a single binary with Postgres as its only hard dependency, and I wanted to keep it that way. So I solved both problems with the database that was already there.
This post covers:
- How the cache works, and how it compares with Redis and Memcached.
- How the queue works, and how it compares with Kafka, RabbitMQ, SQS, Redis and NATS.
- What is happening inside Postgres that makes each technique work, and what it costs.
- How to decide, and the signs that you have outgrown it.
A note on language. The examples are written in Go because that is what the project used. The solution is not a Go solution. The part that does the work is SQL, and it behaves the same from Node.js, Python, Java, Ruby, C# or anything else with a Postgres driver. Near the end there is a section that maps each Go detail to its equivalent elsewhere. If you do want background on the Go pieces, I have written about context and goroutines and channels.
## First, what the traditional tools are for
It helps to name what each dedicated tool is built to do before deciding whether Postgres can stand in for it.
### Key-value stores
A key-value store answers one question very fast: "what is stored under this key?" The two you will meet most are Redis and Memcached.
| Tool | What it is | Known for |
|---|---|---|
| Redis | An in-memory data structure server with optional persistence to disk | Strings, hashes, lists, sets, sorted sets, streams, pub/sub, TTL on every key |
| Memcached | An in-memory cache and nothing else | Simple get and set, multi-threaded, least-recently-used eviction, no disk |
Both keep data in memory, which is why a read is typically answered in well under a millisecond on a local network. Both can be told how much memory they may use, and both evict old entries on their own when that limit is reached.
### Messaging systems
Messaging tools fall into two families that are easy to confuse.
A queue hands each message to one worker. When the worker confirms it is done, the message is gone. This is for distributing work: send this email, resize this image, charge this card.
A log appends every message to an ordered record that is kept for a set time. Many independent readers each track their own position in it and can go back and read again. This is for streaming events: every order placed, every click, every change to a row.
| Tool | Family | What it is |
|---|---|---|
| RabbitMQ | Queue | A broker. Producers publish to exchanges, which route messages into queues by rule |
| Amazon SQS | Queue | A fully managed queue. Nothing to run, pay per request |
| Apache Kafka | Log | A distributed, partitioned, replicated commit log built for very high throughput |
| NATS | Pub/sub and log | A lightweight messaging server. Its JetStream layer adds persistence and replay |
| Redis Streams | Log | An append-only stream type inside Redis, with consumer groups |
| Redis lists, pub/sub | Queue and pub/sub | The simplest options. Lists work as a basic queue. Pub/sub drops messages nobody hears |
Keep those two families in mind. A Postgres table is a good queue. It is a poor log.
## Part one: the cache
### The problem: an expensive answer, asked for constantly
An endpoint is slow because its query is expensive, and it is called far more often than its answer changes. The fix is to compute the answer once, keep it somewhere cheap to read, and serve that copy until it is too old or the underlying data changes.
"Somewhere cheap to read" is usually taken to mean memory, in another process. It does not have to. It only has to be much cheaper than the query it replaces. Looking up one row by its primary key is among the cheapest things a database does.
### The solution: an UNLOGGED table
CREATE UNLOGGED TABLE IF NOT EXISTS cache (
key text PRIMARY KEY,
value jsonb NOT NULL,
expires_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX IF NOT EXISTS idx_cache_expires
ON cache (expires_at) WHERE expires_at IS NOT NULL;The word doing the work is UNLOGGED, and to see why, you need to know what Postgres normally does on a write.
Before Postgres changes a table, it appends a description of the change to the write-ahead log, a sequential file on disk, and makes sure that record is safely stored before telling you the transaction committed. If the server loses power a moment later, it replays the log on restart and nothing committed is lost. Replicas work by receiving and replaying the same log. This is the foundation of Postgres durability, and it is also a real cost: every write is written twice, once to the log and later to the table itself.
An unlogged table opts out. Its changes are not written to the log, so writes are much cheaper. The price has two parts:
- The table is emptied after a crash or an unclean shutdown.
- Its contents are not sent to standby servers, so you cannot read the cache from a replica.
For a cache that is the right trade. Cached data is disposable by definition. Losing it costs one slow request per key.
### Get and set
type Cache struct {
db *pgxpool.Pool
log *slog.Logger
group singleflight.Group
}
// Get returns the value for key, or ok=false on a miss or an expired entry.
// An error is treated as a miss, so a cache failure never breaks a request.
func (c *Cache) Get(ctx context.Context, key string) (json.RawMessage, bool) {
var value json.RawMessage
err := c.db.QueryRow(ctx, `
SELECT value FROM cache
WHERE key = $1 AND (expires_at IS NULL OR expires_at > now())`,
key,
).Scan(&value)
if err != nil {
return nil, false
}
return value, true
}
// Set stores value under key. A ttl of zero or less means no expiry.
func (c *Cache) Set(ctx context.Context, key string, value json.RawMessage, ttl time.Duration) {
var expiresAt *time.Time
if ttl > 0 {
t := time.Now().Add(ttl)
expiresAt = &t
}
ctx, cancel := detach(ctx)
defer cancel()
_, err := c.db.Exec(ctx, `
INSERT INTO cache (key, value, expires_at) VALUES ($1, $2, $3)
ON CONFLICT (key) DO UPDATE
SET value = EXCLUDED.value, expires_at = EXCLUDED.expires_at`,
key, value, expiresAt,
)
if err != nil {
c.log.Warn("cache set failed", "key", key, "err", err)
}
}Two decisions are worth pointing at.
Expiry is checked on read. An expired row is simply not returned, so correctness never depends on a cleanup job having run.
A cache error is a miss, not a failure. If the cache is unhealthy the request falls through to the real query. A cache that can take the API down is worse than no cache.
### Sweeping expired rows
Redis and Memcached remove expired keys for you. Postgres does not know these rows are a cache, so a small loop does it.
// SweepLoop deletes expired entries until ctx is cancelled.
func (c *Cache) SweepLoop(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if _, err := c.db.Exec(ctx, `DELETE FROM cache WHERE expires_at < now()`); err != nil {
c.log.Warn("cache sweep failed", "err", err)
}
}
}
}The partial index on expires_at is there for this query. The sweep is housekeeping only. If it stops, reads stay correct and the table just grows.
### Detaching from the request
detach is three lines, and it fixes a bug that is very hard to see.
// detach returns a context that survives the request being cancelled, but
// still cannot hang.
func detach(ctx context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
}Think about invalidation. A handler updates a row, then deletes the cache key. If the client disconnects between those two steps, the request context is cancelled and the delete never runs. The database has the new value, the cache has the old one, and nothing will correct it until the TTL runs out.
context.WithoutCancel keeps the values of the request context but drops its cancellation. The timeout on top means this bookkeeping still cannot hold a shutdown open.
### Stopping the stampede
The interesting part of a cache is the moment an entry expires. If fifty requests are in flight for the same key, all fifty miss, and all fifty run the expensive query at once. The cache turns a busy moment into a pile-up.
func (c *Cache) GetOrSet(
ctx context.Context,
key string,
ttl time.Duration,
load func() (json.RawMessage, error),
) (json.RawMessage, error) {
if v, ok := c.Get(ctx, key); ok {
return v, nil
}
v, err, _ := c.group.Do(key, func() (any, error) {
// Check again: the request that filled the key may have finished
// while this one waited for its turn.
if v, ok := c.Get(ctx, key); ok {
return v, nil
}
v, err := load()
if err != nil {
return nil, err
}
c.Set(ctx, key, v, ttl)
return v, nil
})
if err != nil {
return nil, err
}
return v.(json.RawMessage), nil
}singleflight makes concurrent callers for the same key share one call to load. The other forty-nine wait and receive the same result. It only covers one process, so three instances can still run the query three times. That is fine. Three is not fifty.
This problem is not specific to Postgres. A Redis cache stampedes in exactly the same way, and needs the same fix.
### Using it in a handler
This is the cache-aside pattern. The handler asks the cache, and only does the real work on a miss.
func (h *Handler) ListCourses(w http.ResponseWriter, r *http.Request) {
filters := parseFilters(r)
// The key is built from the parsed filters, not the raw query string, so
// "?a=1&b=2" and "?b=2&a=1" share one entry.
key := "courses:list:" + filters.CacheKey()
body, err := h.cache.GetOrSet(r.Context(), key, 5*time.Minute, func() (json.RawMessage, error) {
courses, err := h.store.ListCourses(r.Context(), filters)
if err != nil {
return nil, err
}
return json.Marshal(courses)
})
if err != nil {
http.Error(w, "could not load courses", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
w.Write(body)
}Storing the finished JSON has a nice side effect. On a hit there is no decoding and no encoding. The bytes go from the row to the response.
### Invalidating a list
A list endpoint's key has to include its filters, so one write makes a whole family of keys stale: courses:list:page=1, courses:list:page=2&sort=name, and so on. Nobody can enumerate them. So keys are namespaced with a colon and deleted by prefix.
var likeEscaper = strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`)
func (c *Cache) DeletePrefix(ctx context.Context, prefix string) {
ctx, cancel := detach(ctx)
defer cancel()
_, err := c.db.Exec(ctx,
`DELETE FROM cache WHERE key LIKE $1`,
likeEscaper.Replace(prefix)+"%",
)
if err != nil {
c.log.Error("cache invalidation failed, entries may be stale", "prefix", prefix, "err", err)
}
}The escaping matters. In LIKE, an underscore matches any single character. Without escaping, clearing user_roles: would also clear userXroles:. Backslash is replaced first, or the escapes being added would themselves be escaped.
This is one place where Postgres is the easier tool. Deleting by pattern in Redis means walking the keyspace with SCAN, or maintaining your own set of keys per group. Here it is one indexed DELETE.
### Postgres against Redis and Memcached
| Question | Postgres unlogged table | Redis | Memcached |
|---|---|---|---|
| Where the data lives | Disk pages, held in Postgres shared memory | Memory | Memory |
| Speed of a read | A full SQL round trip | Typically well under a millisecond | Typically well under a millisecond |
| Expiry | A column you check and a loop you write | Built in, per key | Built in, per key |
| Bounded memory with eviction | No. The table grows until you delete rows | Yes, with a chosen eviction policy | Yes, least recently used |
| Survives a restart | A clean restart yes, a crash no | Optional, with snapshots or an append-only file | No |
| Readable from replicas | No, unlogged tables are not replicated | Yes | No replication built in |
| Delete a group of keys | One DELETE with LIKE | SCAN and delete, or track keys yourself | Not supported |
| Same transaction as your data | Yes | No | No |
| Richer structures | Anything SQL can express, including JSON | Lists, sets, sorted sets, hashes, streams, and more | Plain values only |
| Extra service to run | None | One | One |
### Advantages of caching in Postgres
- Nothing new to operate. No second service to deploy, patch, monitor, secure, or connect to. One connection string, one backup, one set of credentials.
- Local development is trivial. Anyone who can run the app can run the cache, because it is the app's database.
- Invalidation can join the write. You can delete a cache row in the same transaction that changes the data, so the two cannot disagree. No other option in the table offers this.
- Pattern deletes are easy. One statement clears a whole family of keys.
- You can inspect it with SQL. How many entries, how big, which are oldest, which prefix takes the most space. It is a table.
- The same pool, the same tooling. Metrics, tracing and query logs you already have cover the cache too.
### Disadvantages of caching in Postgres
- It is slower per read. Every hit is a SQL round trip that has to be parsed, planned and executed. For an endpoint that takes 200 milliseconds to compute, that is a huge win. For one that takes two, it may be no win at all.
- It spends database connections. Cache reads compete with real queries for the same pool. Under load, the cache can be the thing that exhausts it.
- There is no eviction. Redis and Memcached drop old entries when memory fills. This table grows until your sweep deletes rows. One bug that writes unbounded keys can fill a disk.
- Updates leave dead rows. Postgres does not overwrite a row in place. Each
UPDATEwrites a new version and leaves the old one for vacuum to clear. A cache with heavy churn creates a lot of that work. - A crash empties it. After an unclean shutdown everything is cold at once, and the database absorbs the full load at the moment it has just restarted.
- It shares fate with the database. A dedicated cache can keep serving when the database is struggling. This one cannot, because it is the database.
## Part two: the queue
### The problem: slow work inside a request
Sending during the request ties the caller's latency, and success, to someone else's SMTP server. A slow relay stalls the handler. A brief outage loses the message unless the user tries again.
Writing a row first and delivering it in the background separates the two. The API answers as soon as the message is safely recorded, and delivery can be retried without anyone resubmitting anything.
Here the queue is the emails table itself. Those rows had to be stored anyway as a record of what was sent, so the queue is three extra columns.
CREATE TABLE IF NOT EXISTS emails (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
recipients text[] NOT NULL,
subject text NOT NULL,
body text NOT NULL DEFAULT '',
status text NOT NULL DEFAULT 'queued'
CHECK (status IN ('queued', 'sending', 'sent', 'failed')),
attempts integer NOT NULL DEFAULT 0,
claimed_at timestamptz,
error text,
sent_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);
-- The worker's claim query: oldest queued rows first.
CREATE INDEX IF NOT EXISTS idx_emails_queue
ON emails (created_at) WHERE status = 'queued';
-- The reclaim sweep, which only looks at rows in flight.
CREATE INDEX IF NOT EXISTS idx_emails_claimed
ON emails (claimed_at) WHERE status = 'sending';Unlike the cache, this table is a normal logged table. A queued email must survive a crash.
Both indexes are partial. The table keeps every email ever sent, but the indexes only hold rows that are queued or in flight, so they stay tiny however much history builds up.
### The biggest advantage: enqueue in the same transaction
This is the single strongest reason to queue in your database, and no external broker can match it.
Say a registration must be saved and a confirmation email must be sent. With a separate broker you have two systems and two writes. Whichever order you choose, one failure leaves them disagreeing:
- Save first, publish second. If the publish fails, the user is registered and never hears about it.
- Publish first, save second. If the save fails, the user gets a confirmation for a registration that does not exist.
This is known as the dual-write problem. The standard cure is the outbox pattern: write the message into a table in the same transaction as the data, then have something forward it to the broker. When the queue is already a table, you get the outbox for free and skip the forwarding.
func (s *Service) Register(ctx context.Context, in Registration) error {
tx, err := s.db.Begin(ctx)
if err != nil {
return err
}
// A no-op once the transaction has been committed.
defer tx.Rollback(ctx)
if _, err := tx.Exec(ctx,
`INSERT INTO registrations (user_id, course_id) VALUES ($1, $2)`,
in.UserID, in.CourseID,
); err != nil {
return err
}
if _, err := tx.Exec(ctx,
`INSERT INTO emails (recipients, subject, body) VALUES ($1, $2, $3)`,
[]string{in.Email}, "You are registered", renderConfirmation(in),
); err != nil {
return err
}
// Both rows exist, or neither does.
return tx.Commit(ctx)
}There is no window in which one happened and the other did not.
### The naive queue, and why it breaks
A table of pending work is the easy part. The hard part is two workers. The obvious way to take a message is to read one and then mark it:
SELECT id FROM emails WHERE status = 'queued' ORDER BY created_at LIMIT 1;
-- then, with the id that came back:
UPDATE emails SET status = 'sending' WHERE id = $1;With one worker this is fine. With two, it sends duplicates, because nothing stops both from reading the same row before either has marked it.
| Step | Worker A | Worker B |
|---|---|---|
| 1 | Reads the oldest queued row, number 7 | |
| 2 | Reads the oldest queued row, also number 7 | |
| 3 | Marks 7 as sending | |
| 4 | Marks 7 as sending, which also succeeds | |
| 5 | Sends the email | Sends the same email |
The textbook fix is to lock the row as you read it, by adding FOR UPDATE to the select. That removes the duplicate and creates a different problem. Worker B now waits at step 2 until A's transaction ends. When it does, B looks at row 7 again, finds it is no longer queued, and comes back with nothing, even though other messages are sitting right behind it. The workers end up taking turns. Adding more of them adds no speed.
What you want is for B to ignore the row A is holding and take the next one.
### The solution: claiming work with SKIP LOCKED
This one statement is what makes Postgres usable as a queue.
UPDATE emails
SET status = 'sending', attempts = attempts + 1, claimed_at = now()
WHERE id IN (
SELECT id FROM emails
WHERE status = 'queued' AND attempts < $1
ORDER BY created_at
LIMIT $2
FOR UPDATE SKIP LOCKED
)
RETURNING id, recipients, subject, body, attempts;FOR UPDATE locks the rows the inner query selects. Without SKIP LOCKED, a second worker running the same query would wait for the first one's locks, then find the rows already claimed. With it, the second worker steps over locked rows and takes the next ones. Several workers can run at once and each message goes to exactly one of them.
The claim and the status change happen in one statement, so there is no gap in which a row is selected but not yet marked.
Notice what this does not do. It does not hold a transaction open while the email is being sent. The claim commits at once and the row is marked sending. A long transaction around slow network work would pin a connection and hold back vacuum for the whole table, which is the classic way a Postgres queue gets into trouble.
### The worker loop
func (s *Service) RunWorker(ctx context.Context, cfg WorkerConfig) {
ticker := time.NewTicker(cfg.Interval)
defer ticker.Stop()
for {
// Recover anything a dead worker left behind before taking new work.
if n, err := s.reclaimStale(ctx); err != nil {
s.log.Error("reclaim failed", "err", err)
} else if n > 0 {
s.log.Warn("requeued abandoned email", "count", n)
}
// Keep going while full batches come back, so a backlog drains at the
// speed of the relay and not one batch per tick.
for {
claimed, err := s.deliverBatch(ctx, cfg.BatchSize)
if err != nil {
s.log.Error("sweep failed", "err", err)
break
}
if claimed < cfg.BatchSize {
break
}
}
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}It polls. The interval only decides how quickly a new message is noticed on an idle system, because a backlog is drained in a tight inner loop. For email, fifteen seconds is nothing.
### Waking the worker at once with LISTEN and NOTIFY
Polling was enough for my case, so I did not build this. If you need a message picked up immediately, Postgres can signal the worker.
The producer adds one statement to its transaction. The notification is only delivered if the transaction commits.
SELECT pg_notify('email_queued', '');The worker holds one connection open and waits on it, with the polling interval as a fallback.
func (s *Service) waitForWork(ctx context.Context, conn *pgx.Conn, interval time.Duration) {
wait, cancel := context.WithTimeout(ctx, interval)
defer cancel()
// Returns when a notification arrives or the interval passes. Either way
// the caller sweeps next, so an error here needs no handling.
_, _ = conn.WaitForNotification(wait)
}Keep the fallback. A notification is not stored. If no worker is listening when it is sent, it is gone. That is fine here because the row is still in the table, and the next poll finds it. The notification is a hint to hurry, never the record of the work.
### When a worker dies mid-delivery
A row is sending while a worker holds it. If that worker is killed, the row stays sending for ever, and no claim query will ever pick it up. So each pass starts by returning old claims to the queue.
UPDATE emails
SET status = 'queued', claimed_at = NULL
WHERE status = 'sending'
AND claimed_at < now() - make_interval(secs => $1);The timeout is a real trade-off, so choose it on purpose. A message requeued while its first delivery was only slow gets sent twice. A message never requeued is never sent. A duplicate notice is a much smaller harm than a verification email that never arrives, so I set the timeout to five minutes, well above the thirty second SMTP timeout, and accept the remaining risk.
If you have used Amazon SQS this will look familiar. It is the same idea as the SQS visibility timeout: a received message is hidden from other consumers for a while, and reappears if nobody deletes it in time.
### Recording the outcome during shutdown
func (s *Service) deliverBatch(ctx context.Context, n int) (int, error) {
batch, err := s.claim(ctx, n, maxAttempts)
if err != nil {
return 0, err
}
for _, msg := range batch {
sendErr := s.mailer.Send(ctx, msg)
// The outcome is written on its own context. Shutdown cancels ctx, and
// a message delivered a moment before that would otherwise fail to
// record its own success and be sent again on the next boot.
outcome, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
if sendErr == nil {
s.markSent(outcome, msg.ID)
} else {
// Back to 'queued' for another try, or 'failed' once the attempts
// are spent, so a typo in an address cannot spin the queue for ever.
s.markFailed(outcome, msg.ID, sendErr, msg.Attempts >= maxAttempts)
}
cancel()
}
return len(batch), nil
}This is the same WithoutCancel idea as the cache, for the same reason. The side effect has already happened, so the record of it must not depend on a context that may be cancelled at that moment.
The attempts cap matters as much as the retry. Without it, one permanently undeliverable address is retried for ever and every sweep spends time on it.
### What delivery guarantee is this?
Every messaging system makes one of three promises. It is worth knowing which one you have.
| Guarantee | Meaning | What goes wrong | Typical cause |
|---|---|---|---|
| At most once | A message is delivered zero or one times | Messages can be lost | Acknowledging before the work is done |
| At least once | A message is delivered one or more times | Messages can be duplicated | Acknowledging after the work, with retries |
| Exactly once | The effect of a message happens once | Hard to achieve, and never free | Needs idempotent work or transactions end to end |
The queue in this post is at least once. A crash between sending the email and recording the outcome means the email is sent again after the reclaim timeout.
No system gives you exactly once delivery to an outside service such as an SMTP relay, because the broker cannot make the relay's send and its own acknowledgement one atomic step. What systems call exactly once is really at least once delivery combined with work that is safe to repeat. So if a duplicate would be dangerous, make the handler idempotent, for example by recording a unique key for each message and refusing to process the same key twice.
### The same ideas under different names
Dedicated brokers have names for everything this table does by hand. The mapping is close to one to one.
| Concept | In a broker | In this table |
|---|---|---|
| Take a message | Receive, consume, pop | The claim query with SKIP LOCKED |
| Confirm it is done | Acknowledge, delete, commit offset | Set status to sent |
| Hide it while it is processed | Visibility timeout (SQS), unacknowledged (RabbitMQ) | status = 'sending' with claimed_at |
| Return it if the worker died | Redelivery, XAUTOCLAIM in Redis Streams | The reclaim sweep |
| Park it after too many failures | Dead letter queue or exchange | status = 'failed' once attempts are spent |
| Several workers, one message each | Competing consumers, a consumer group | Several workers running the same claim query |
| How far behind are we | Queue depth, consumer lag | count(*) where status is queued |
### Postgres against the brokers
| Question | Postgres table | RabbitMQ | Amazon SQS | Apache Kafka | Redis Streams | NATS JetStream |
|---|---|---|---|---|---|---|
| Family | Queue | Queue | Queue | Log | Log | Log |
| Default guarantee | At least once | At least once with acks | At least once | At least once | At least once with groups | At least once |
| Ordering | Roughly by insert time | Per queue | Best effort, strict with FIFO queues | Strict within a partition | Strict within a stream | Strict within a stream |
| Replay old messages | Only if you keep rows | Not on classic queues | No | Yes, for the retention period | Yes | Yes |
| Many independent readers | Build it yourself | Yes, through exchanges | Through SNS fan-out | Yes, consumer groups | Yes, consumer groups | Yes, consumers |
| Routing rules | Any SQL WHERE | Rich: direct, topic, fanout | None | By topic and partition key | By stream | By subject, with wildcards |
| Same transaction as your data | Yes | No | No | No | No | No |
| Throughput ceiling | Modest | High | High, and managed for you | Very high | High | Very high |
| What you operate | Nothing new | A broker, often a cluster | Nothing, it is a service | A cluster of brokers | Redis | A NATS server or cluster |
A few of those rows deserve a sentence.
Ordering. With one worker, this table is processed in insert order. With several, two workers can finish out of order, and a retried message jumps behind newer ones. Email does not care. A stream of account balance changes does. If order matters, Kafka's rule is the one to copy: everything that must stay in order shares a key and is handled by one consumer.
Replay. Kafka keeps messages after they are read, so a new service can start from the beginning and catch up on history. That is the defining feature of a log, and it is why Kafka is used for event sourcing, analytics pipelines and feeding several systems from one stream. A queue table has nothing like it unless you add a position per reader, at which point you are building a small Kafka badly.
Fan-out. In this table a message has one status, so it belongs to one consumer. Delivering the same event to five independent services needs five rows or a second table of per-consumer positions. RabbitMQ exchanges, Kafka consumer groups and NATS subjects do this out of the box.
Throughput. Each message here costs an insert, an update to claim it, and an update to finish it, all through the write-ahead log, and each leaves dead row versions for vacuum. That is comfortable at tens or hundreds of messages a second on ordinary hardware. A system built as an append-only log does far less work per message and scales much further.
### Advantages of queuing in Postgres
- Atomic with your data. A message is enqueued if and only if the transaction commits. This removes a whole class of bugs.
- Nothing new to operate. No broker cluster, no client library with its own failure modes, no extra credentials or network rules.
- Durable by default. Rows go through the same write-ahead log, replication and backups as everything else you store.
- Fully inspectable. Finding a stuck message is a
SELECT. Retrying a failed one is anUPDATE. Most brokers make both of those harder than they should be. - The history is free. The table is already an audit trail of what was sent, to whom, when, and with what error.
- Flexible selection. Priority, scheduled delivery and per-tenant fairness are an extra column and a different
ORDER BYorWHERE. - Safe with many workers.
SKIP LOCKEDgives competing consumers without any coordination service.
### Disadvantages of queuing in Postgres
- A lower ceiling. It will not carry the message volumes Kafka or NATS are built for.
- Table churn. Every status change is a new row version. Busy queues need vacuum to keep up, and long-running transactions elsewhere in the database can stop it from doing so.
- Polling adds latency. Without
LISTENandNOTIFY, a message waits up to one interval before it is seen. - No replay, no fan-out, no routing. These have to be built, and each one is real work.
- It uses your database's resources. Workers hold connections, and a backlog adds load to the system that serves your users.
- You own the edge cases. Reclaim timeouts, retry caps, back-off and poison messages are code you write and test. A broker ships them.
- It only suits one database. Services that do not share the database cannot share the queue.
### Why a queue table needs looking after
The warnings above about dead rows and vacuum come from how Postgres stores changes, and it is worth understanding because it explains every piece of advice that follows.
Postgres never overwrites a row in place. An UPDATE writes a complete new version of the row and marks the old version as dead. A DELETE only marks the row dead. The dead versions stay on disk until a background process called vacuum confirms that no running transaction could still need to see them, and then frees the space for reuse. This design is called multi-version concurrency control, and it is what lets readers and writers work at the same time without blocking each other.
A queue is close to a worst case for it. Each message is inserted once and updated at least twice, so a queue that has handled a million messages has produced at least two million dead row versions. Normally vacuum clears them as fast as they appear and nobody notices.
The trouble starts when vacuum cannot keep up, and the usual reason is one long-running transaction somewhere in the database. Vacuum may not remove any row version that an open transaction might still be able to see. A report that runs for an hour, or a transaction left open by a forgotten session, therefore pins every dead row created since it began. The table and its indexes swell, the claim query has to step over more and more dead entries to find a live one, and the queue slows down for a reason that has nothing to do with the queue.
That is the background to three small habits.
Watch it. One query gives you depth and age, which are the two numbers that matter.
SELECT status,
count(*) AS messages,
now() - min(created_at) AS oldest
FROM emails
WHERE status IN ('queued', 'sending', 'failed')
GROUP BY status;A growing queued count means workers are not keeping up. An old sending row means the reclaim sweep is not running. A rising failed count means something downstream is rejecting messages.
Keep transactions short. Claim, commit, do the work, then record the outcome. Never hold a transaction open across a network call.
Decide what happens to finished rows. I keep them because they are the audit trail, and the partial indexes mean they cost nothing at claim time. If you do not need the history, delete or archive completed rows on a schedule so the table does not grow without bound.
### You do not have to write this yourself
I wrote it by hand because the need was small and the rows already existed. If you want more features, several mature projects use exactly this approach, across most languages:
- River is a job queue for Go built on Postgres, with retries, scheduling, unique jobs and a web interface.
- pgmq is a Postgres extension that gives you an SQS-style queue as SQL functions.
- Oban for Elixir, pg-boss for Node.js, Que for Ruby and Solid Queue for Rails all follow the same design.
The fact that so many ecosystems have arrived at the same place is a good sign that the idea is sound.
## The same trick, a third time
Once the pattern is in place it keeps paying off. Rate limiting across several API instances needs a shared counter, which is usually another reason to add Redis, where you would pair INCR with EXPIRE. Here it is one upsert on another unlogged table.
INSERT INTO rate_limits (key, window_start, hits, expires_at)
VALUES ($1, $2, 1, $3)
ON CONFLICT (key, window_start)
DO UPDATE SET hits = rate_limits.hits + 1
RETURNING hits;func (l *Limiter) Allow(ctx context.Context, key string) bool {
// Truncating to the window means every instance agrees on the bucket
// boundary without talking to each other.
start := time.Now().Truncate(l.window)
var hits int
err := l.db.QueryRow(ctx, incrementSQL, key, start, start.Add(l.window)).Scan(&hits)
if err != nil {
// Fail open. A limiter is a safeguard, not a reason to take the API
// down when the database hiccups.
return true
}
return hits <= l.limit
}Concurrent requests, including ones on other instances, queue up on the row, so two of them can never read the same stale count.
This is a fixed window, so a caller can make up to twice the limit across a window boundary: a full allowance at the end of one minute and another at the start of the next. That is the accepted price for a single round trip and no extra bookkeeping. A sliding window is more accurate and costs more per request.
## None of this is specific to Go
Everything above that matters is either SQL or an idea. The SQL runs unchanged from any language. The ideas each have a direct equivalent wherever you work.
| The idea | In the Go examples | Elsewhere |
|---|---|---|
| The cache table, upsert, expiry check and prefix delete | SQL sent through pgx | The same SQL through node-postgres, psycopg, JDBC, Npgsql or any other driver |
| Concurrent misses on one key share a single load | singleflight.Group | Keep the in-flight promise, future or task in a map keyed by cache key, and hand that same one to every caller |
| Bookkeeping finishes even if the client disconnects | context.WithoutCancel with a timeout | Do not tie the write to the request's abort signal or cancellation token, and give it its own short timeout |
| The claim query and the reclaim sweep | SQL sent through pgx | The same SQL. FOR UPDATE SKIP LOCKED is a feature of Postgres, not of the client |
| A worker that sweeps on an interval | A goroutine with a time.Ticker | A background task, a timer loop, a separate worker process, or a scheduled job |
| Enqueue with the data, atomically | pool.Begin, two inserts, Commit | Any transaction API. Both inserts go inside one transaction |
| Stop cleanly on shutdown | A cancelled context ends the loop | Handle the termination signal, stop claiming, and let in-flight work record its outcome |
Two of those are worth a closer look because they are the ones most often got wrong in any language.
Sharing one load per key. The mistake is to cache only the finished value. Between the first miss and the value arriving, every other caller also sees a miss. The fix is to cache the work in progress as well. The first caller starts the load and stores the pending result. Everyone who arrives before it finishes is handed that same pending result and waits on it. When it resolves, it is removed from the map. That is all singleflight does, and it is a dozen lines in most languages.
Finishing after the client has gone. Many frameworks cancel work when the client disconnects, which is right for reads and wrong for cleanup. Look at how your framework propagates cancellation, and make sure cache invalidation and "this message was delivered" writes are not on that path.
## How to decide
### Postgres is the right call when
- You already run Postgres and would be adding a service only for this.
- Message volume is modest, and a delay of a few seconds is acceptable.
- The work is naturally tied to a database write, so enqueueing in the same transaction has real value.
- The queued items are data you would store anyway.
- The cached endpoints are slow to compute, so a SQL round trip is still a large saving.
- The team is small and every extra piece of infrastructure has a real cost in attention.
### Reach for a dedicated tool when
| If you need | Use |
|---|---|
| Reads in well under a millisecond, thousands of times a second | Redis or Memcached |
| A cache with a hard memory limit and automatic eviction | Redis or Memcached |
| Sorted sets, counters, leaderboards, sessions, distributed locks | Redis |
| The simplest possible cache with nothing else attached | Memcached |
| Complex routing between many producers and consumers | RabbitMQ |
| A queue with no servers to manage at all | Amazon SQS |
| An event stream that several services read independently | Apache Kafka |
| To replay history, or rebuild state from past events | Apache Kafka |
| Very high throughput across many machines | Apache Kafka or NATS JetStream |
| Low-latency messaging between many small services | NATS |
| Streams, and you already run Redis | Redis Streams |
### Signs you have outgrown it
- Cache reads show up near the top of your slowest or most frequent queries.
- The connection pool runs out under load and cache or worker traffic is a large share of it.
- The queue table keeps growing in size even though the number of live messages is steady.
- You find yourself adding a per-consumer position column, or copying messages so a second service can read them.
- Other teams or services want to consume the same messages.
None of these is a failure. They are the point at which the extra infrastructure starts paying for itself, and moving is straightforward because the interface is small. The cache is four methods. The queue is enqueue, claim, and record an outcome. Both can be reimplemented on Redis or a broker without touching the code that calls them.
## Summary
- An unlogged table is a perfectly good cache when the thing being cached is slow to compute. It has no eviction, so sweep it.
- Share one load per key to stop a stampede when an entry expires, whatever the cache is built on. In Go that is
singleflight. - Detach invalidation and outcome writes from request cancellation, or a disconnect will leave stale state behind.
FOR UPDATE SKIP LOCKEDturns a table into a queue that several workers can share safely, from any language.- Postgres never updates a row in place, so keep transactions short and let vacuum do its job.
- Enqueueing in the same transaction as the data is the advantage no broker can offer.
- The result is at least once. Make handlers safe to repeat.
- Postgres is a good queue and a poor log. If you need replay, fan-out or very high throughput, use the tool built for it.
One database you already run, back up and understand is a real feature. The whole thing here is two tables, a few hundred lines of application code, and nothing new to operate. Start there, measure, and add the dedicated tool on the day the numbers ask for it.
Published on October 7, 2026
38 min read
Found an Issue!
Find an issue with this post? Think you could clarify, update or add something? All my posts are available to edit on Github. Any fix, little or small, is appreciated!
Edit on GitHubLast updated on
