Skip to main content

A schema-first FastStream framework with durable retries, dead-letter handling, and optional ClickHouse support.

Project description

easy-faststream

easy-faststream is a schema-first wrapper around FastStream for RabbitMQ.

Application developers define a Pydantic schema and consumer function. The framework handles validation, acknowledgements, durable retries, dead-letter routing, and optional idempotent ClickHouse insertion.

Installation

Install the RabbitMQ framework:

pip install easy-faststream

Install it with ClickHouse support:

pip install "easy-faststream[clickhouse]"

Basic consumer

from datetime import datetime
from uuid import UUID

from pydantic import BaseModel

from easy_faststream import MessageContext, StreamApp


class TripRequested(BaseModel):
    order_id: UUID
    passenger_id: UUID | None = None
    service_type: str
    event_at: datetime


stream = StreamApp.from_env()


@stream.consumer(
    event="passapp.trip.requested",
    schema=TripRequested,
    queue="p_q.passapp.trip.requested",
    exchange="ex.passapp.event.trip",
    routing_key="passapp.trip.requested",
)
async def consume_trip(
    event: TripRequested,
    context: MessageContext,
) -> None:
    print(event.order_id, context.message_id)


app = stream.app

Run the application:

python -m faststream run app:app

RabbitMQ configuration

EASY_STREAM_RABBITMQ_URL=amqp://guest:guest@localhost:5672/
EASY_STREAM_APP_NAME=trip-consumer

EASY_STREAM_DEFAULT_EXCHANGE=easy.events
EASY_STREAM_RETRY_EXCHANGE=easy.events.retry
EASY_STREAM_DLQ_EXCHANGE=easy.events.dead

EASY_STREAM_RETRY_QUEUE_SUFFIX=.retry
EASY_STREAM_DLQ_QUEUE_SUFFIX=.dead

EASY_STREAM_MAX_RETRIES=3
EASY_STREAM_RETRY_DELAY_SECONDS=1
EASY_STREAM_RETRY_BACKOFF=2

EASY_STREAM_GRACEFUL_TIMEOUT=30

With the settings above, retry delays are 1, 2, and 4 seconds.

Retries use durable RabbitMQ TTL queues. Retry messages therefore survive consumer shutdowns and application restarts.

ClickHouse sink

Configure the ClickHouse connection:

EASY_STREAM_CLICKHOUSE_HOST=localhost
EASY_STREAM_CLICKHOUSE_PORT=8123
EASY_STREAM_CLICKHOUSE_USERNAME=default
EASY_STREAM_CLICKHOUSE_PASSWORD=
EASY_STREAM_CLICKHOUSE_DATABASE=default
EASY_STREAM_CLICKHOUSE_SECURE=false

Attach a sink to a consumer:

from easy_faststream import ClickHouseSink, StreamApp


stream = StreamApp.from_env()

sink = ClickHouseSink(
    table="bronze.trip_requested",
    idempotency_key="order_id",
)


@stream.consumer(
    event="passapp.trip.requested",
    schema=TripRequested,
    queue="p_q.passapp.trip.requested",
    exchange="ex.passapp.event.trip",
    routing_key="passapp.trip.requested",
    sink=sink,
)
async def consume_trip(event: TripRequested) -> None:
    print(f"Processing order {event.order_id}")


app = stream.app

The processing order is:

Validate schema
    → Run handler
    → Insert into ClickHouse
    → ACK RabbitMQ message

If ClickHouse insertion fails, the RabbitMQ message is not considered successfully processed.

ClickHouse idempotency

The sink creates a stable SHA-256 insert_deduplication_token using:

  • ClickHouse table
  • RabbitMQ routing key
  • Configured business key, such as order_id

Publishing the same business event with different RabbitMQ message IDs therefore uses the same ClickHouse token.

The destination table must use a MergeTree family engine with an appropriate deduplication window. ClickHouse deduplication is limited by that window and is not a permanent unique-key constraint.

Permanent and transient failures

Transient failures, such as connection timeouts, enter the durable retry flow.

Permanent ClickHouse errors are sent directly to the DLQ without unnecessary retries. Examples include:

  • Missing table or database
  • Unknown column
  • Type mismatch
  • Invalid input
  • SQL syntax error

Failure behavior

  • Invalid schema: publish the original payload and validation error to the DLQ.
  • Processing error: use durable retry queues, then publish to the DLQ.
  • Transient ClickHouse error: use the durable retry flow.
  • Permanent ClickHouse error: publish directly to the DLQ.
  • NonRetryableError: skip retries and publish directly to the DLQ.
  • Retry publication failure: NACK and requeue the original message.
  • DLQ publication failure: NACK and requeue the original message.
  • Successful processing: ACK the original message.
  • Successful DLQ publication: ACK the original message.

Development

Install development dependencies:

python -m pip install -e ".[dev]"

Run checks:

python -m ruff check src tests examples
python -m pytest
python -m build
python -m twine check dist/*

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

easy_faststream-0.3.0.tar.gz (5.3 MB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

easy_faststream-0.3.0-py3-none-any.whl (13.0 kB view details)

Uploaded Python 3

File details

Details for the file easy_faststream-0.3.0.tar.gz.

File metadata

  • Download URL: easy_faststream-0.3.0.tar.gz
  • Upload date:
  • Size: 5.3 MB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.14.6

File hashes

Hashes for easy_faststream-0.3.0.tar.gz
Algorithm Hash digest
SHA256 675c9e0fa5fe277de25b1e43756dc20d5ac3e0a166abd5ae1424ac2b027f0e2b
MD5 f416c4e74162c0c04abc2272b9a06d43
BLAKE2b-256 bd1160e8d7d18b4f36f51fe24785fa7b8307dd19845201ee9b9813cd98227028

See more details on using hashes here.

File details

Details for the file easy_faststream-0.3.0-py3-none-any.whl.

File metadata

File hashes

Hashes for easy_faststream-0.3.0-py3-none-any.whl
Algorithm Hash digest
SHA256 30cae1be3f9a035b0f277adfc46f3f3b409b601ce95cd5612d1fd3d89a32babf
MD5 c6417ecab26a35130f698215eca287b9
BLAKE2b-256 39516a9894107d7e830db9d7c16150b4c9b24e3f8132814682c6674cfe758509

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page