Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 28 additions & 29 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,20 @@
# `pgoutbox` - a transactional outbox for `pgx`

`pgoutbox` implements a simple [transactional outbox](https://microservices.io/patterns/data/transactional-outbox.html) for [`pgx`](https://github.com/jackc/pgx). New messages can be added to a Postgres table using `AddMessages` and can be flushed to a destination via `ProcessMessages`.
[![Go Reference](https://pkg.go.dev/badge/github.com/hatchet-dev/pgoutbox.svg)](https://pkg.go.dev/github.com/hatchet-dev/pgoutbox)

`pgoutbox` implements a simple [transactional outbox](https://microservices.io/patterns/data/transactional-outbox.html) for [`pgx`](https://github.com/jackc/pgx). New messages can be added to a Postgres table within a transaction using `AddMessages` and can be flushed to a destination via `ProcessMessages`.

## Why?

While working on [Hatchet](https://github.com/hatchet-dev/hatchet) we needed a reliable and performant way to durably persist messages over a message boundary. In particular, we needed:

- **Batched reads and writes.** `AddMessages` inserts a batch of messages in a single transaction, and `ProcessMessages` locks a batch, hands the whole batch to one `Flush` call, and deletes it in the same transaction (see [Atomic flush and delete](#atomic-flush-and-delete) and [Benchmarks](#benchmarks)).
- **Exclusive consumers with leasing semantics.** Exactly one instance across a fleet owns a topic under a renewing lease, and a standby takes over within seconds if the holder goes away (see [Exclusive consumers](#exclusive-consumers)).
- **Support for publishing across hundreds of thousands of topics.** Topics are plain strings that don't need to be declared up front, and they're tracked in a table rather than by a poller or worker pool per topic (see [Multiple topics and flushers](#multiple-topics-and-flushers) and [Message expiration](#message-expiration)).

Without these particular requirements, a library like [River](https://github.com/riverqueue/river) would otherwise have been a good fit. `pgoutbox` is deliberately an outbox rather than a job queue: there are no retries, scheduling, priorities, or job history, and messages are deleted as soon as they're flushed.

## Example usage

Here's an example of flushing messages on `topic1` by simply printing them to the console:

Expand Down Expand Up @@ -44,7 +58,7 @@ if err != nil {
}
```

## Schema
## Schema migrations

By default, `NewOutbox` runs migrations and creates an outbox table in the schema `outbox.messages`. This can be overwritten via:

Expand Down Expand Up @@ -97,7 +111,7 @@ outbox.ProcessMessages(ctx, "shipments")

## Atomic flush and delete

If your flusher writes to Postgres itself (e.g. into a relay table), use the transaction exposed by the `FlushContext` passed to `Flush`. It is the same transaction `ProcessMessages` uses to lock and delete messages, so your writes and the outbox delete commit or roll back together:
If your flusher writes to Postgres itself (e.g. into a separate table), use the transaction exposed by the `FlushContext` passed to `Flush`. It is the same transaction `ProcessMessages` uses to lock and delete messages, so your writes and the outbox delete commit or roll back together:

```go
type relayFlusher struct{}
Expand All @@ -117,7 +131,7 @@ func (f *relayFlusher) Flush(ctx pgoutbox.FlushContext, msgs []*sqlc.Message) er

## Continuous processing with `Subscribe`

Instead of calling `ProcessMessages` yourself, `Subscribe` runs it in a loop: it drains the topic, then waits until either the poll interval elapses or a new-message notification arrives (see below), and drains again. It blocks until its context is cancelled:
Instead of calling `ProcessMessages` yourself, `Subscribe` drains the topic automatically by listening for a poll or a new-message notification. It blocks until its context is cancelled:

```go
go func() {
Expand All @@ -131,7 +145,7 @@ go func() {
}()
```

Processing errors don't kill the loop — they're logged to the `WithLogger` logger and retried on the next wake-up. `Subscribe` returns an error immediately only if no flusher is registered for the topic or the subscription itself can't be established.
If processing fails, messages are logged to the `WithLogger` logger and retried again on the next wake-up. `Subscribe` returns an error immediately only if no flusher is registered for the topic or the subscription itself can't be established.

### Waking on new messages with `LISTEN`/`NOTIFY`

Expand All @@ -146,9 +160,7 @@ if err != nil {
outbox, err := pgoutbox.NewOutbox(ctx, pool, pgoutbox.WithPubSub(ps))
```

With a `PubSub` attached, `AddMessages` publishes a notification for each staged topic and `Subscribe` wakes on it instead of waiting out the poll interval. The pg-backed `PubSub` publishes *inside the `AddMessages` transaction*, so the notification is delivered exactly when the insert commits — and never for a transaction that rolls back. Postgres deduplicates identical notifications within a transaction, so any number of `AddMessages` calls for a topic in one transaction cost a single wake-up.

A notification only makes sense once the staging transaction has committed, and a `PubSub` that doesn't implement `TxPublisher` has no way to defer a publish to commit time. For those transports, pass a `Notifier` to `AddMessages` and fire it after a successful commit:
You can swap your own pub/sub implementation as well. Note that a notification only makes sense once the inserting transaction has committed, so a `PubSub` that doesn't implement `TxPublisher` needs to defer notifying until after commit. In these cases, pass a `Notifier` to `AddMessages` and fire it after a successful commit:

```go
var notifier pgoutbox.Notifier
Expand All @@ -167,14 +179,9 @@ if err := tx.Commit(ctx); err != nil {
notifier.Notify(ctx) // wakes the subscribers of both topics
```

The `Notifier` accumulates one notification per `AddMessages` call it is passed to, so a single one can serve a whole transaction. With a `TxPublisher` transport like the pg-backed `PubSub`, `Notify` is a no-op (the notification already rode the transaction), so the pattern is transport-agnostic. Skipping it never loses messages — subscribers just fall back to the poll interval — and publish failures inside `Notify` are logged, not returned, for the same reason.
Delivery is best-effort by design; if notifications are lost, `pgoutbox` falls back to polling.

Delivery is best-effort by design: if a notification is lost (for example while the listener reconnects), polling picks the messages up within one poll interval. The listener occupies a single dedicated connection (hijacked out of the pool so it doesn't consume a pool slot) no matter how many topics are subscribed.

Two details worth knowing:

- All notifications travel over one NOTIFY channel, `pgoutbox_pubsub` by default. Two outboxes sharing a database (e.g. different schemas) should use distinct channels via `pgoutbox.WithNotifyChannel("my_channel")` to avoid waking each other's subscribers.
- `PubSub` is an interface, so you can bring your own transport (e.g. Redis, NATS) instead of `LISTEN`/`NOTIFY`. If your implementation also implements `TxPublisher` (detected once, at `NewOutbox`), notifications are published transactionally as described above; otherwise use `WithNotifier` and fire `Notify` after commit to publish best-effort.
The default LISTEN/NOTIFY channel is `pgoutbox_pubsub`, which can be configured via `pgoutbox.WithNotifyChannel("my_channel")` if necessary.

## Message expiration

Expand All @@ -194,7 +201,7 @@ if err != nil {
}
```

The outbox always runs a background scanner goroutine that polls the `topics` table, but maintenance loops are only launched for topics that actually have an expiration configured. Topics don't need to be declared at startup: the library tracks every topic that receives a message in a `topics` table (via a Postgres trigger) and applies the default expiration automatically.
The outbox always runs a background scanner goroutine that polls the `topics` table, but maintenance loops are only launched for topics that actually have an expiration configured. Topics don't need to be declared at startup: the library tracks every topic that receives a message in a `topics` table and applies the default expiration automatically.

Multiple outbox instances (e.g. replicas of the same service) coordinate cleanup using a per-topic maintenance lease, so only one instance runs the delete at a time.

Expand All @@ -207,11 +214,9 @@ outbox, err := pgoutbox.NewOutbox(ctx, pool,
)
```

Lease competition — another instance winning the cleanup race — is not logged.

## Exclusive consumers

By default, any number of `ProcessMessages` callers can drain a topic concurrently (each call grabs a non-overlapping batch using `FOR UPDATE SKIP LOCKED`). Use `AcquireTopic` when you need exactly one active consumer at a time.
By default, any number of `ProcessMessages` callers can drain a topic concurrently (using `FOR UPDATE SKIP LOCKED`). Use `AcquireTopic` or `pgoutbox.WithExclusive()` passed to `Subscribe` when you need exactly one active consumer at a time.

`AcquireTopic` blocks until this instance holds the exclusive lease for the topic, then returns. A background goroutine automatically renews the lease until the context is cancelled, at which point the lease expires and another instance can take over:

Expand Down Expand Up @@ -240,9 +245,7 @@ An instance that never acquired the lease (or whose lease has expired) receives

When the holder's context is cancelled, the lease expires naturally (within the lease duration, 30 s by default) and another instance's `AcquireTopic` call unblocks.

### Exclusive subscribers

`Subscribe` composes with exclusive consumers: pass `WithExclusive()` and it manages the lease for you.
You can also pass `pgoutbox.WithExclusive()` to `Subscribe`:

```go
// Exactly one instance across the fleet drains "orders" at a time; the rest
Expand All @@ -256,19 +259,17 @@ With `WithExclusive()`, `Subscribe`:
2. re-acquires it automatically if the lease is ever lost mid-subscribe (for example, a heartbeat lapse during a database blip); and
3. releases it on return, so a waiting instance takes over immediately instead of waiting out the lease's grace period.

Several instances calling `Subscribe(..., WithExclusive())` on the same topic therefore form a failover group: one active consumer, the rest hot standbys. Nothing is missed during a handoff — the new holder's first drain pass covers any backlog that accumulated while the lease changed hands.

Without `WithExclusive()`, subscribing to a topic whose exclusive lease is held elsewhere doesn't fail — every pass errors (visible via `WithLogger`) and is retried, so the subscriber sits idle until the lease frees up or is acquired. For an exclusive topic, either call `AcquireTopic` before `Subscribe` or pass `WithExclusive()`.
Several instances calling `Subscribe(..., WithExclusive())` on the same topic therefore form a failover group: one active consumer, the rest hot standbys. Without `WithExclusive()`, subscribing to a topic whose exclusive lease is held elsewhere doesn't fail — every pass errors (visible via `WithLogger`) and is retried, so the subscriber sits idle until the lease frees up or is acquired. For an exclusive topic, either call `AcquireTopic` before `Subscribe` or pass `WithExclusive()`.

## Benchmarks

You can run benchmarks locally; for example, to write and flush 100k messages, you can run:
On a local Macbook with an M3 Max core, `pgoutbox` reaches `223817 msgs/sec` with batching support. You can run benchmarks locally; for example, to write and flush 100k messages, you can run:

```
go test -bench=. -benchtime=100000x
```

`BenchmarkOutbox_WriteAndPublishThroughput` drains each topic with a busy-polling `ProcessMessages` loop; `BenchmarkOutbox_SubscribeThroughput` drains each topic with a `Subscribe` call woken by `LISTEN`/`NOTIFY` (its poll interval is set far above the benchmark runtime, so throughput is carried entirely by notifications — and each producer commit pays the in-transaction `pg_notify`). Both run a matrix of 1 and 10 topics (messages spread round-robin, one consumer per topic) and producer batch sizes of 1, 10, and 100 messages per `AddMessages` transaction. On a local Macbook with an M3 Max core:
Full results from a local run:

```
$ go test -bench=. -benchtime=100000x
Expand Down Expand Up @@ -301,5 +302,3 @@ BenchmarkOutbox_SubscribeThroughput/TxFlush/topics=10/batch=10-14 1
BenchmarkOutbox_SubscribeThroughput/Flush/topics=10/batch=100-14 100000 5352 ns/op 186840 msgs/sec
BenchmarkOutbox_SubscribeThroughput/TxFlush/topics=10/batch=100-14 100000 39021 ns/op 25627 msgs/sec
```

Batching producer writes is the single biggest lever: staging 100 messages per transaction reaches ~150-220k msgs/sec on the cheap-flusher path, roughly 20× the one-message-per-transaction rate. The `batch=1` cells run long enough to be sensitive to background database noise (checkpoints, autovacuum), so expect swings between runs — the 1082 msgs/sec outlier above measures ~3300 msgs/sec in isolation.
48 changes: 48 additions & 0 deletions doc.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
// Package pgoutbox implements a simple transactional outbox for pgx.
//
// New messages are added to a Postgres table within a transaction using
// [Outbox.AddMessages] and are flushed to a destination via
// [Outbox.ProcessMessages]. Because the insert rides the caller's own
// transaction, a message is only ever visible to consumers if the business
// write it belongs to committed.
//
// Messages are grouped by topic, and each topic is drained by the [Flusher]
// registered for it with [Outbox.AddFlusher]. ProcessMessages locks a batch of
// messages, hands the whole batch to a single Flush call, and deletes the
// messages in the same transaction if the flush succeeds. A Flusher that
// writes to Postgres itself can use [FlushContext.Tx] so its writes and the
// outbox delete commit or roll back together.
//
// [Outbox.Subscribe] drains a topic continuously, waking on a poll interval
// or, when a [PubSub] such as [NewPGPubSub] is attached via [WithPubSub], the
// moment new messages commit. [Outbox.AcquireTopic] and [WithExclusive] give
// a topic exactly one active consumer across a fleet, backed by a renewing
// lease with automatic failover. [WithTopicExpiration] and
// [WithDefaultExpiration] delete old messages in the background.
//
// A minimal setup looks like:
//
// type printFlusher struct{}
//
// func (printFlusher) Flush(_ pgoutbox.FlushContext, msgs []*sqlc.Message) error {
// for _, m := range msgs {
// fmt.Printf("flushed id=%d topic=%s payload=%s\n", m.ID, m.Topic, string(m.Payload))
// }
// return nil
// }
//
// outbox, err := pgoutbox.NewOutbox(ctx, pool)
// if err != nil {
// return err
// }
// outbox.AddFlusher("orders", printFlusher{})
//
// // within a transaction
// err = outbox.AddMessages(ctx, tx, "orders", []pgoutbox.MessageOpts{{Payload: payload}})
//
// // after the transaction commits
// _, err = outbox.ProcessMessages(ctx, "orders")
//
// See the README at https://github.com/hatchet-dev/pgoutbox for a longer
// walkthrough of each feature.
package pgoutbox
13 changes: 6 additions & 7 deletions migrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,14 +30,13 @@ func validateSchemaName(schema string) error {
return nil
}

// Migrate runs the embedded pgoutbox migrations against the given pool.
// It is the explicit alternative to NewOutbox's auto-migration: callers
// that want to control when DDL runs (separate startup phase, release
// pipeline, etc.) should construct the outbox with WithAutoMigrate(false)
// and invoke Migrate themselves.
// Migrate runs the embedded pgoutbox migrations against pool, creating the
// schema if needed. NewOutbox runs them automatically by default. To run them
// yourself instead, for example as part of a separate release step, construct
// the outbox with WithAutoMigrate(false) and call Migrate explicitly.
//
// Only WithSchema is consulted from opts; other options are accepted for
// API symmetry but ignored.
// Only WithSchema is consulted from opts. Other options are accepted so the
// same option list can be passed to NewOutbox, but they are ignored here.
func Migrate(ctx context.Context, pool *pgxpool.Pool, opts ...OutboxOpt) error {
o := defaultOpts()
for _, f := range opts {
Expand Down
4 changes: 4 additions & 0 deletions no_op_flusher.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,16 @@ package pgoutbox

import "github.com/hatchet-dev/pgoutbox/sqlc"

// NopFlusher is a Flusher that discards every message it is given. It is useful
// in tests, and for draining a topic whose messages are no longer needed.
type NopFlusher struct{}

// NewNopFlusher returns a NopFlusher.
func NewNopFlusher() *NopFlusher {
return &NopFlusher{}
}

// Flush discards msgs and returns nil.
func (f *NopFlusher) Flush(_ FlushContext, _ []*sqlc.Message) error {
return nil
}
Loading
Loading