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
28 changes: 14 additions & 14 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,11 @@
[![Ruff](https://img.shields.io/endpoint?url=https://raw.githubusercontent.com/astral-sh/ruff/main/assets/badge/v2.json)](https://github.com/astral-sh/ruff)
[![ty](https://img.shields.io/endpoint?url=https://raw.githubusercontent.com/astral-sh/ty/main/assets/badge/v0.json)](https://github.com/astral-sh/ty)

`faststream-outbox` is a [FastStream](https://faststream.ag2.ai) broker integration for the **transactional outbox pattern** — a Postgres table is the message queue.
`faststream-outbox` is a [FastStream](https://faststream.ag2.ai) broker integration for the transactional outbox pattern, with a Postgres table as the message queue.

A producer writes a domain entity and an outbox row in the *same* SQLAlchemy transaction by calling `broker.publish(body, queue=..., session=session)`. A separate subscriber polls the table and relays each row to a real message bus (Kafka, RabbitMQ, NATS, Redis…) with a single decorator — or processes the rows in-place if you don't have a downstream broker.
A producer writes a domain entity and an outbox row in the *same* SQLAlchemy transaction by calling `broker.publish(body, queue=..., session=session)`. A separate subscriber polls the table and relays each row to a real message bus (Kafka, RabbitMQ, NATS, Redis…) with a single decorator, or processes the rows in place if you don't have a downstream broker.

## Quickstart — outbox relay to Kafka
## Quickstart: outbox relay to Kafka

Write the outbox row in your domain transaction; relay rows to Kafka with a stacked decorator.

Expand Down Expand Up @@ -65,9 +65,9 @@ async with session_factory() as session, session.begin():

The same one-decorator pattern works for RabbitMQ, NATS, Redis, and Confluent. See the [relay tutorial](https://faststream-outbox.modern-python.org/usage/relay/) for the FastAPI lifecycle, header propagation, router shapes, and the at-least-once contract.

## Quickstart — standalone outbox queue
## Quickstart: standalone outbox queue

If you don't have a downstream broker, the same broker can process outbox rows in-place — the table *is* the queue.
If you don't have a downstream broker, the same broker can process outbox rows in place, and the table *is* the queue.

```python
from sqlalchemy import MetaData
Expand Down Expand Up @@ -97,25 +97,25 @@ async with session_factory() as session, session.begin():

## How it works

A subscriber owns two async loops: a **fetch** loop claims available rows via a single CTE (`SELECT … FOR UPDATE SKIP LOCKED → UPDATE acquired_token=:uuid, acquired_at=now() RETURNING *`), and `max_workers` **worker** loops dispatch to the handler. On success, `DELETE WHERE id=:id AND acquired_token=:token`; on failure, the retry strategy schedules another attempt or terminally drops the row. Terminal failures `DELETE` by default; pass `dlq_table=make_dlq_table(metadata)` to atomically archive them into a sibling audit table instead — see [Dead-letter queue](https://faststream-outbox.modern-python.org/usage/dlq/).
A subscriber owns two async loops: a fetch loop claims available rows via a single CTE (`SELECT … FOR UPDATE SKIP LOCKED → UPDATE acquired_token=:uuid, acquired_at=now() RETURNING *`), and `max_workers` worker loops dispatch to the handler. On success, `DELETE WHERE id=:id AND acquired_token=:token`; on failure, the retry strategy schedules another attempt or terminally drops the row. Terminal failures `DELETE` by default; pass `dlq_table=make_dlq_table(metadata)` to atomically archive them into a sibling audit table instead (see [Dead-letter queue](https://faststream-outbox.modern-python.org/usage/dlq/)).

The `acquired_token` is the load-bearing invariant: a slow handler whose lease expired and was re-claimed by another worker finds its terminal `DELETE` to be a no-op (the token no longer matches), preventing it from clobbering the new lease holder.
The `acquired_token` is the load-bearing invariant. If a slow handler's lease expired and another worker re-claimed the row, the slow handler's terminal `DELETE` is a no-op because the token no longer matches, so it cannot clobber the new lease holder.

With the `asyncpg` driver, the fetch loop also `LISTEN`s on `outbox_<table>` and `publish` emits `pg_notify(...)`, so idle dispatch latency is ~10ms instead of up to `max_fetch_interval`.

See [How it works](https://faststream-outbox.modern-python.org/introduction/how-it-works/) for the full architecture.

## Optional extras

- `faststream-outbox[asyncpg]` — asyncpg driver (enables `LISTEN/NOTIFY` for ~10ms idle dispatch)
- `faststream-outbox[fastapi]` — FastAPI integration via `faststream_outbox.fastapi.OutboxRouter`
- `faststream-outbox[validate]` — Alembic for `broker.validate_schema()`
- `faststream-outbox[prometheus]` — Prometheus metrics adapter
- `faststream-outbox[opentelemetry]` — OpenTelemetry metrics adapter
- `faststream-outbox[asyncpg]`: asyncpg driver (enables `LISTEN/NOTIFY` for ~10ms idle dispatch)
- `faststream-outbox[fastapi]`: FastAPI integration via `faststream_outbox.fastapi.OutboxRouter`
- `faststream-outbox[validate]`: Alembic for `broker.validate_schema()`
- `faststream-outbox[prometheus]`: Prometheus metrics adapter
- `faststream-outbox[opentelemetry]`: OpenTelemetry metrics adapter

## Acknowledgements

The architecture of this package is heavily informed by Arseniy Popov's [PR #2704](https://github.com/ag2ai/faststream/pull/2704) (`feat: add sqla broker`) on upstream FastStream — the FastStream broker/registrator/subscriber wiring, the `SELECT … FOR UPDATE SKIP LOCKED` fetch-and-claim CTE, the retry strategy hierarchy, and the in-transaction publish contract all originate from there. This package is a Postgres-only reimplementation that diverges in storage model (lease tokens instead of an explicit state column, archive table is opt-in), loop structure (two loops instead of four), wake-up mechanism (`LISTEN/NOTIFY`), and adds timer mechanics. Credit for the original design belongs to Arseniy.
The architecture of this package is heavily informed by Arseniy Popov's [PR #2704](https://github.com/ag2ai/faststream/pull/2704) (`feat: add sqla broker`) on upstream FastStream. The FastStream broker/registrator/subscriber wiring, the `SELECT … FOR UPDATE SKIP LOCKED` fetch-and-claim CTE, the retry strategy hierarchy, and the in-transaction publish contract all originate from there. This package is a Postgres-only reimplementation. It differs in storage model (lease tokens in place of an explicit state column, and an opt-in archive table), loop structure (two loops where the PR has four), and wake-up mechanism (`LISTEN/NOTIFY`), and it adds timer mechanics. Credit for the original design belongs to Arseniy.

## 📚 [Documentation](https://faststream-outbox.modern-python.org)

Expand All @@ -126,4 +126,4 @@ The architecture of this package is heavily informed by Arseniy Popov's [PR #270
## Part of `modern-python`

Browse the full list of templates and libraries in
[`modern-python`](https://github.com/modern-python) — see the org profile for the categorized index.
[`modern-python`](https://github.com/modern-python); the org profile has the categorized index.
Binary file modified docs/assets/social-card.png
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
87 changes: 42 additions & 45 deletions docs/concepts/comparison.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,75 +16,74 @@ reader can lift the answer without reading the discussion.

## vs. writing your own outbox table and worker

A bespoke outbox is the most common starting point — the pattern itself is
A bespoke outbox is the most common starting point. The pattern itself is
straightforward, and an MVP can ship in an afternoon. What `faststream-outbox`
buys you, in concrete terms, is the pile of pieces that turn that MVP into a
production system:
gives you is the pile of pieces that turn that MVP into a production system:

- per-row **lease tokens** with a load-bearing invariant that any new
- per-row lease tokens with a load-bearing invariant that any new
fetch / terminal path must preserve;
- the **partial-index design** the fetch CTE depends on (without those, the
- the partial-index design the fetch CTE depends on (without those, the
disjunctive WHERE clause falls back to seq-scan as the table grows);
- the **fetch-and-claim CTE shape** with `FOR UPDATE SKIP LOCKED` that
- the fetch-and-claim CTE shape with `FOR UPDATE SKIP LOCKED` that
reclaims both unleased rows and expired leases without a separate reaper;
- the **retry-strategy hierarchy** (`ExponentialRetry`, `NoRetry`, …)
- the retry-strategy hierarchy (`ExponentialRetry`, `NoRetry`, …)
enforcing `max_attempts` and `max_total_delay_seconds` uniformly;
- **`validate_schema()`** via Alembic's `autogenerate.compare_metadata`;
- **drain semantics** on stop: new fetches stop while in-flight handlers
- `validate_schema()` via Alembic's `autogenerate.compare_metadata`;
- drain semantics on stop: new fetches stop while in-flight handlers
finish, and subscribers drain concurrently;
- **`LISTEN/NOTIFY`** short-circuit on top of polling, with NOTIFY
- a `LISTEN/NOTIFY` short-circuit on top of polling, with NOTIFY
suppression on future-dated rows and `timer_id` conflict no-ops;
- **`timer_id` dedup** via a partial unique index plus
- `timer_id` dedup via a partial unique index plus
`on_conflict_do_nothing`;
- the **DLQ atomicity CTE** that rolls back the DELETE when the DLQ insert
- the DLQ atomicity CTE that rolls back the DELETE when the DLQ insert
fails.

None of these is hard individually. Together they make up most of the work
past an MVP.

You also pick up the [Subscriber](../usage/subscriber.md),
You also get the [Subscriber](../usage/subscriber.md),
[Publisher](../usage/publisher.md), [Dead-letter queue](../usage/dlq.md), and
[Observability](../usage/observability.md) reference pages — written to
[Observability](../usage/observability.md) reference pages, written to
the level you would otherwise have to write yourself.

**TL;DR.** Build it yourself if you only ever need the MVP shape. Use
Build it yourself if you only ever need the MVP shape. Use
`faststream-outbox` if you expect the system to live for a couple of years.

## vs. CDC (Debezium, logical replication)

Change-data capture sits one layer below outbox: instead of writing rows
to an outbox table, you read your write-ahead log directly. The producer
code is unchanged — any write to the underlying tables becomes an event.
code is unchanged, and any write to the underlying tables becomes an event.
Debezium and similar tools have spent the last decade hardening this path
for Postgres, MySQL, and others; the operator playbook is well-known.

CDC wins when you already need WAL-level capture for analytics or
reverse-ETL anyway, when you want to capture writes from services you do
not control (i.e. not all your producers are FastStream apps), or when
the polling overhead of an outbox is unacceptable. CDC also wins when
the events you care about are derivable from row state — "an order
exists with status='paid'" rather than "an
the events you care about are derivable from row state, such as "an order
exists with status='paid'" as opposed to "an
`OrderPaid` event was published."

`faststream-outbox` wins when you control the producer code (so the
outbox row is cheap to write inline with the domain write), when you
need **handler-level retry, DLQ, and scheduled-delivery semantics
inline** (CDC pushes those concerns to a separate consumer layer), and
when the **async-Python logical-replication tooling gap** is too thin
to lean on: there is no async-native logical-decoding client comparable
need handler-level retry, DLQ, and scheduled-delivery semantics
inline (CDC pushes those concerns to a separate consumer layer), and
when the async-Python logical-replication tooling is too thin
to lean on. There is no async-native logical-decoding client comparable
to Debezium's JVM connectors, so a Python CDC path means either running
the JVM stack alongside your app or driving `pg_recvlogical` / a thin
`psycopg` replication-protocol wrapper yourself. That gap is why this
project does not recommend CDC as the default path.

**TL;DR.** Pick CDC when you already need WAL capture or have producers
Pick CDC when you already need WAL capture or have producers
outside your control. Pick this when you own the producer and want
retry/DLQ/timers in-process.

## vs. Kafka transactions (or RabbitMQ publisher confirms)

Atomic `DB-write + bus-publish` is also achievable on a real bus, just
not for free. Kafka transactions plus two-phase commit, or an
Atomic `DB-write + bus-publish` is also achievable on a real bus, at a
cost. Kafka transactions plus two-phase commit, or an
idempotent-producer pattern combined with an inbox table on the consumer
side, can give you the same end-to-end at-least-once guarantee without a
DB-backed outbox.
Expand All @@ -100,12 +99,12 @@ to preserve, because the transactional boundary belongs to two different
systems.

`faststream-outbox` covers a focused subset of a real bus's delivery
surface — one-process producer, one Postgres table — while adding
surface (one-process producer, one Postgres table) while adding
outbox-native features a bare bus lacks (`cancel_timer`, `timer_id`
producer-side dedup, scheduled delivery), at the price of being
Postgres-only and polling-based.

**TL;DR.** Kafka transactions / Rabbit confirms win at scale where the
Kafka transactions / Rabbit confirms win at scale where the
bus is already running. `faststream-outbox` wins when Postgres is your
only durable store.

Expand All @@ -117,56 +116,54 @@ across listener disconnect: a NOTIFY emitted while the listener's
connection is dead, or during a reconnect, is silently dropped. There is
no replay, no persistence, no retry.

`faststream-outbox` keeps the **outbox row** as the durability boundary
`faststream-outbox` keeps the outbox row as the durability boundary
and uses NOTIFY only as a wake-up short-circuit on top of polling. If
the NOTIFY is lost — listener reconnecting, or `LISTEN` setup failed at
startup — the subscriber still finds the row on its next poll cycle. The worst case
the NOTIFY is lost (listener reconnecting, or `LISTEN` setup failed at
startup), the subscriber still finds the row on its next poll cycle. The worst case
is one `max_fetch_interval` of idle latency (default 10 seconds), not
data loss.

**TL;DR.** Raw `LISTEN/NOTIFY` is a wake-up, not a delivery guarantee.
Raw `LISTEN/NOTIFY` is a wake-up, not a delivery guarantee.
Use the outbox row for durability and let NOTIFY shave idle latency.

## vs. Celery (or RQ, Dramatiq) with a DB backend

Celery and friends are *task queues* — you submit "go do this thing" and
Celery and friends are *task queues*: you submit "go do this thing" and
a worker picks it up later. `faststream-outbox` is *message routing* with
FastStream's subscriber/publisher semantics — the row is an event tied
FastStream's subscriber/publisher semantics: the row is an event tied
to a domain write, and the handler is the consumer for events on that
queue.

The two abstractions overlap, but the right one depends on what you are
modelling. Celery wins for ad-hoc background jobs initiated from
arbitrary points in your app (request handlers, admin commands, cron),
where the relationship to a database transaction is incidental. Use
`faststream-outbox` when you want **at-least-once dispatch of events
that must commit atomically with a domain write**, and prefer
`faststream-outbox` when you want at-least-once dispatch of events
that must commit atomically with a domain write, and prefer
FastStream's `@broker.subscriber` model over Celery's task decorator.

The two can also coexist — Celery for fire-and-forget background jobs,
`faststream-outbox` for the transactional event tier. They are not in
direct competition for the same problem.
The two can also coexist: Celery for fire-and-forget background jobs,
`faststream-outbox` for the transactional event tier.

**TL;DR.** Celery for ad-hoc background jobs. `faststream-outbox` for
Celery for ad-hoc background jobs. `faststream-outbox` for
events tied to DB transactions.

## vs. FastStream + `KafkaBroker` / `RabbitBroker` directly

If you have no domain write to atomically commit alongside the bus
publish, drop the outbox entirely — use the foreign broker directly via
publish, drop the outbox entirely and use the foreign broker directly via
FastStream's native `KafkaBroker`, `RabbitBroker`, `NatsBroker`, etc.
You skip the polling overhead and the Postgres dependency; you keep the
same `@broker.subscriber` ergonomics.

The interesting case is **both at once**: domain code writes to Postgres
The interesting case is both at once: domain code writes to Postgres
*and* needs the event to reach Kafka. That is the canonical
transactional-outbox shape, and it composes the two: the outbox row
captures the event in the domain transaction; a
[Relay](../usage/relay.md) subscriber forwards it to Kafka with the
at-least-once contract preserved end to end. Don't pick between
`faststream-outbox` and a real bus — use both, with the outbox as the
durability boundary in front of the bus.
at-least-once contract preserved end to end. Use both, with the outbox
as the durability boundary in front of the bus.

**TL;DR.** No DB write to commit with? Use the foreign broker directly.
No DB write to commit with? Use the foreign broker directly.
Need atomicity with a DB write? Use this *plus* the foreign broker via
[Relay](../usage/relay.md).
Loading
Loading