NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
PyPI · #3609 most downloaded on PyPI
A thin, typed Python client for Kafka, RabbitMQ, NATS, Redis and MQTT — with AsyncAPI docs and in-memory testing.
Last release 11 days ago
23 Sep 2026
Ships fairly regularly
a new release about every 4 weeks
Nearly every release is documented
notes for 60 of the last 60 stable releases
Nothing withdrawn
no release was ever pulled
3 years old
126 releases · first in 2023
One column per quarter.
ci: stop a macro error at the docs build, and drop diff-cover's deprecated flag by @Lancetnik in #3218
Full Changelog: 0.7.6...0.7.7
feat(confluent): add Topic schema to configure topic creation by @vyhuholl in #3026
address field carries the address (#2357) by @Lancetnik in #3038assert_called_once_with for subscribers and publishers by @Lancetnik in #3110assert_called_with and assert_any_call beside assert_called_once_with by @Lancetnik in #3111docs serve by @SarthakB11 in #2874subscriber() is typed by max_workers, so a decorated handler stays typed by @Lancetnik in #3123no_ack has no effect with a consumer group by @Griger10 in #3129TopicPartition is FastStream's own type by @Lancetnik in #3087Address.literal and Address.describe by @Lancetnik in #3104Full Changelog: 0.7.5...0.7.6
Update Release Notes for 0.7.4 by @Lancetnik in #2998
realign_keys method in faststream/response/response.py by @ApusBerliozi in #2988Full Changelog: 0.7.4...0.7.5
fix(TestClient-redis): [#2963] Fixed groups for TestRedisBroker by @ApusBerliozi in https://github.com/ag2ai/faststream/pull/2965
ws_connection_headers and reconnect_to_server_handler params" by @IvanKirpichnikov in https://github.com/ag2ai/faststream/pull/2971Full Changelog: https://github.com/ag2ai/faststream/compare/0.7.3...0.7.4
ws_connection_headers and reconnect_to_server_handler params" by @IvanKirpichnikov in #2971Full Changelog: 0.7.3...0.7.4
Feature/deprecate fastapi plugin by @IvanKirpichnikov in #2944
typing.Required hints for sentinel redis by @IvanKirpichnikov in #2955Full Changelog: 0.7.2...0.7.3
feat(redis): add Redis Sentinel support by @PAzter1101 in https://github.com/ag2ai/faststream/pull/2895
Full Changelog: https://github.com/ag2ai/faststream/compare/0.7.1...0.7.2
Full Changelog: 0.7.1...0.7.2
TestBroker.__aenter__ was typed to return Broker | list[Broker]. That union is wrong for both usage shapes: mypy rejects .publish() on the single-brok
TestBroker.aenter was typed to return Broker | list[Broker]. That union is wrong for both usage shapes: mypy rejects .publish() on the single-broker result (the list arm has no such method) and rejects unpacking the multi-broker result (the Broker arm is not iterable).
# Before — both lines fail under `mypy`:
async with TestKafkaBroker(KafkaBroker()) as br:
await br.publish(None, "test")
# error: Item "list[KafkaBroker]" of "KafkaBroker | list[KafkaBroker]" has no attribute "publish" [union-attr]
async with TestKafkaBroker(KafkaBroker(), KafkaBroker()) as (br1, br2):
# error: "KafkaBroker" object is not iterable [misc]
...
# After — mypy infers the precise type:
async with TestKafkaBroker(KafkaBroker()) as br:
reveal_type(br) # KafkaBroker
await br.publish(None, "test")
async with TestKafkaBroker(KafkaBroker(), KafkaBroker()) as (br1, br2):
reveal_type(br1) # tuple[KafkaBroker, ...] -> KafkaBroker
await br1.publish(None, "test")
await br2.publish(None, "test")
Full Changelog: https://github.com/ag2ai/faststream/compare/0.7.0...0.7.1
TestBroker.aenter was typed to return Broker | list[Broker]. That union is wrong for both usage shapes: mypy rejects .publish() on the single-broker result (the list arm has no such method) and rejects unpacking the multi-broker result (the Broker arm is not iterable).
# Before — both lines fail under `mypy`:
async with TestKafkaBroker(KafkaBroker()) as br:
await br.publish(None, "test")
# error: Item "list[KafkaBroker]" of "KafkaBroker | list[KafkaBroker]" has no attribute "publish" [union-attr]
async with TestKafkaBroker(KafkaBroker(), KafkaBroker()) as (br1, br2):
# error: "KafkaBroker" object is not iterable [misc]
...
# After — mypy infers the precise type:
async with TestKafkaBroker(KafkaBroker()) as br:
reveal_type(br) # KafkaBroker
await br.publish(None, "test")
async with TestKafkaBroker(KafkaBroker(), KafkaBroker()) as (br1, br2):
reveal_type(br1) # tuple[KafkaBroker, ...] -> KafkaBroker
await br1.publish(None, "test")
await br2.publish(None, "test")Full Changelog: 0.7.0...0.7.1
The following APIs that were deprecated in earlier 0.x releases have been fully removed in 0.7.0:
FastStream now includes a full-featured MQTT broker, installable via pip install faststream[mqtt]. It supports wildcard topic filters, path parameter capture via Path(), QoS levels, per-subscriber ack_policy, and AsyncAPI schema generation.
from faststream import FastStream, Path
from faststream.mqtt import MQTTBroker, MQTTMessage, QoS
broker = MQTTBroker("localhost:1883")
app = FastStream(broker)
@broker.subscriber(
"sensors/{device_id}/temperature",
qos=QoS.AT_LEAST_ONCE,
)
async def on_temperature(body: str, device_id: Annotated[str, Path()]) -> None:
print(device_id, body)
@app.after_startup
async def publish_demo() -> None:
await broker.publish(21.5, "sensors/room1/temperature", qos=QoS.AT_LEAST_ONCE)
A single FastStream application can now run multiple brokers at the same time. Pass all the brokers directly to the FastStream constructor — each keeps its own subscribers and publishers, and the app starts and stops all of them together. A common use case is bridging two systems: consume from one broker and re-publish to another.
from faststream import FastStream
from faststream.kafka import KafkaBroker
from faststream.nats import NatsBroker
kafka_broker = KafkaBroker("localhost:9092")
nats_broker = NatsBroker("nats://localhost:4222")
app = FastStream(kafka_broker, nats_broker)
@kafka_broker.subscriber("incoming")
@nats_broker.publisher("outgoing")
async def from_kafka(msg: str) -> str:
# Bridge the message from Kafka to NATS
return msg
@nats_broker.subscriber("outgoing")
async def from_nats(msg: str) -> None:
print(f"Received from NATS: {msg}")
FastStream's Redis broker now has a dedicated RedisClusterBroker that connects to a Redis Cluster with automatic node discovery. It is a drop-in replacement for RedisBroker — just change the class name and point it at any cluster node.
from faststream import FastStream
from faststream.redis import RedisClusterBroker
# A single URL is enough — the cluster auto-discovers all remaining nodes
broker = RedisClusterBroker("redis://node1:7000")
app = FastStream(broker)
@broker.subscriber("events")
async def handle_event(msg: str) -> None:
print(f"Received: {msg}")
@app.after_startup
async def publish_event() -> None:
await broker.publish("hello from cluster", "events")
The AsyncAPIRoute class (used in ASGI hosting) has had two parameters renamed:
| Before | After | Notes |
|---|---|---|
try_it_out=False |
try_it_out_path=None |
Disabling try-it-out now uses None instead of False |
try_it_out_url="..." |
try_it_out_path="..." |
Parameter renamed for clarity |
# Before
AsyncAPIRoute("/docs/asyncapi", try_it_out=False)
AsyncAPIRoute("/docs/asyncapi", try_it_out_url="https://api.example.com/asyncapi/try")
# After
AsyncAPIRoute("/docs/asyncapi", try_it_out_path=None)
AsyncAPIRoute("/docs/asyncapi", try_it_out_path="https://api.example.com/asyncapi/try")
Additionally, a new asyncapi_json_path parameter was added (defaults to <path>.json) and its position in the signature changed — use keyword arguments to avoid surprises.
durable=True is now the default (PR #2892)RabbitQueue and RabbitExchange now default to durable=True (previously False). This aligns with RabbitMQ 4.3+ which disables transient non-exclusive queues by default.
Impact: if you already have a transient (non-durable) queue or exchange of the same name declared on your broker, re-declaration will raise a PRECONDITION_FAILED mismatch error. To opt out, pass durable=False explicitly:
from faststream.rabbit import RabbitQueue
# To keep the old transient behavior:
queue = RabbitQueue("my-queue", durable=False)
The following APIs that were deprecated in earlier 0.x releases have been fully removed in 0.7.0:
ack_first, no_ack and related subscriber options — replaced by ack_policy=AckPolicy.*RedisJSONMessageParser — removed. All Redis services must now use the binary message format.broker.close() — removed. Use broker.stop() instead.Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.7...0.7.0
FastStream now includes a full-featured MQTT broker, installable via pip install faststream[mqtt]. It supports wildcard topic filters, path parameter capture via Path(), QoS levels, per-subscriber ack_policy, and AsyncAPI schema generation.
from faststream import FastStream, Path
from faststream.mqtt import MQTTBroker, MQTTMessage, QoS
broker = MQTTBroker("localhost:1883")
app = FastStream(broker)
@broker.subscriber(
"sensors/{device_id}/temperature",
qos=QoS.AT_LEAST_ONCE,
)
async def on_temperature(body: str, device_id: Annotated[str, Path()]) -> None:
print(device_id, body)
@app.after_startup
async def publish_demo() -> None:
await broker.publish(21.5, "sensors/room1/temperature", qos=QoS.AT_LEAST_ONCE)A single FastStream application can now run multiple brokers at the same time. Pass all the brokers directly to the FastStream constructor — each keeps its own subscribers and publishers, and the app starts and stops all of them together. A common use case is bridging two systems: consume from one broker and re-publish to another.
from faststream import FastStream
from faststream.kafka import KafkaBroker
from faststream.nats import NatsBroker
kafka_broker = KafkaBroker("localhost:9092")
nats_broker = NatsBroker("nats://localhost:4222")
app = FastStream(kafka_broker, nats_broker)
@kafka_broker.subscriber("incoming")
@nats_broker.publisher("outgoing")
async def from_kafka(msg: str) -> str:
# Bridge the message from Kafka to NATS
return msg
@nats_broker.subscriber("outgoing")
async def from_nats(msg: str) -> None:
print(f"Received from NATS: {msg}")FastStream's Redis broker now has a dedicated RedisClusterBroker that connects to a Redis Cluster with automatic node discovery. It is a drop-in replacement for RedisBroker — just change the class name and point it at any cluster node.
from faststream import FastStream
from faststream.redis import RedisClusterBroker
# A single URL is enough — the cluster auto-discovers all remaining nodes
broker = RedisClusterBroker("redis://node1:7000")
app = FastStream(broker)
@broker.subscriber("events")
async def handle_event(msg: str) -> None:
print(f"Received: {msg}")
@app.after_startup
async def publish_event() -> None:
await broker.publish("hello from cluster", "events")The AsyncAPIRoute class (used in ASGI hosting) has had two parameters renamed:
| Before | After | Notes |
|---|---|---|
try_it_out=False |
try_it_out_path=None |
Disabling try-it-out now uses None instead of False |
try_it_out_url="..." |
try_it_out_path="..." |
Parameter renamed for clarity |
# Before
AsyncAPIRoute("/docs/asyncapi", try_it_out=False)
AsyncAPIRoute("/docs/asyncapi", try_it_out_url="https://api.example.com/asyncapi/try")
# After
AsyncAPIRoute("/docs/asyncapi", try_it_out_path=None)
AsyncAPIRoute("/docs/asyncapi", try_it_out_path="https://api.example.com/asyncapi/try")Additionally, a new asyncapi_json_path parameter was added (defaults to <path>.json) and its position in the signature changed — use keyword arguments to avoid surprises.
durable=True is now the default (PR #2892)RabbitQueue and RabbitExchange now default to durable=True (previously False). This aligns with RabbitMQ 4.3+ which disables transient non-exclusive queues by default.
Impact: if you already have a transient (non-durable) queue or exchange of the same name declared on your broker, re-declaration will raise a PRECONDITION_FAILED mismatch error. To opt out, pass durable=False explicitly:
from faststream.rabbit import RabbitQueue
# To keep the old transient behavior:
queue = RabbitQueue("my-queue", durable=False)The following APIs that were deprecated in earlier 0.x releases have been fully removed in 0.7.0:
ack_first, no_ack and related subscriber options — replaced by ack_policy=AckPolicy.*RedisJSONMessageParser — removed. All Redis services must now use the binary message format.broker.close() — removed. Use broker.stop() instead.Full Changelog: 0.6.7...0.7.0
fix: corrected security parsing for mqtt broker by @lemmehoop in https://github.com/ag2ai/faststream/pull/2832
Full Changelog: https://github.com/ag2ai/faststream/compare/0.7.0rc0...0.7.0rc1
Full Changelog: 0.7.0rc0...0.7.0rc1
ack_policy now replaces several deprecated options
Just two main changes:
from faststream.mqtt import MQTTBroker (thanks @borisalekseev)broker.close removed, use broker.stop insteadYou install the release manually
pip install "faststream[mqtt]==0.7.0rc0"
# or
uv add --pre "faststream[mqtt]==0.7.0rc0"
We will release a stable version as soon as we test MQTTBroker with production services (in a few weeks).
Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.7...0.7.0rc0
Just two main changes:
from faststream.mqtt import MQTTBroker (thanks @borisalekseev)broker.close removed, use broker.stop insteadYou install the release manually
pip install "faststream[mqtt]==0.7.0rc0"
# or
uv add --pre "faststream[mqtt]==0.7.0rc0"We will release a stable version as soon as we test MQTTBroker with production services (in a few weeks).
Full Changelog: 0.6.7...0.7.0rc0
The main feature of this release is the Try It Out feature for your Async API documentation!
The main feature of this release is the Try It Out feature for your Async API documentation!
Now you can test your developing application directly from the web, just like Swagger for HTTP. It supports in-memory publication to test a subscriber and real broker publication to verify behavior in real scenarios.
<img width="1467" height="640" alt="Снимок экрана 2026-03-01 в 11 16 10" src="https://github.com/user-attachments/assets/4320e674-24d5-4ead-9820-4bb979e340e7" />
Full updates:
Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.6...0.6.7
The main feature of this release is the Try It Out feature for your Async API documentation!
Now you can test your developing application directly from the web, just like Swagger for HTTP. It supports in-memory publication to test a subscriber and real broker publication to verify behavior in real scenarios.
Full updates:
Full Changelog: 0.6.6...0.6.7
Add support for aiokafka 0.13 by @dolfinus in https://github.com/ag2ai/faststream/pull/2754
F811 by @chirizxc in https://github.com/ag2ai/faststream/pull/2737Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.5...0.6.6
F811 by @chirizxc in #2737Full Changelog: 0.6.5...0.6.6
fix: FastAPI 0.128 compatibility
ServiceUnavailableError for nats by @swelborn in https://github.com/ag2ai/faststream/pull/2720Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.4...0.6.5
feat: Enables message keys for batch publishing by @ozeranskii in https://github.com/ag2ai/faststream/pull/2586
get_one method of stream … by @powersemmi in https://github.com/ag2ai/faststream/pull/2667Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.3...0.6.4
Fix annotation for group_instance_id parameter by @gandhis1 in https://github.com/ag2ai/faststream/pull/2606
min_idle_time in Redis StreamSub and XAUTOCLAIM by @powersemmi in https://github.com/ag2ai/faststream/pull/2607Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.2...0.6.3
Asgi request validation error and docs by @borisalekseev in https://github.com/ag2ai/faststream/pull/2525
xack and xdel methods in Redis testing setup. by @powersemmi in https://github.com/ag2ai/faststream/pull/2599Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.1...0.6.2
feat: add --loop option to run command by @dimastbk in https://github.com/ag2ai/faststream/pull/2572
Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.0...0.6.1
We tried our best to minimize breaking changes, but unfortunately, some aspects were simply not working well. Therefore, we decided to break them in o…
FastStream 0.6 is a significant technical release that aimed to address many of the current project design issues and unlock further improvements on the path to version 1.0.0. We tried our best to minimize breaking changes, but unfortunately, some aspects were simply not working well. Therefore, we decided to break them in order to move forward.
This release includes:
The primary goal of this release is to unlock the path towards further features. Therefore, we are pleased to announce that after this release, we plan to work on MQTT #956 and SQS #794 support and move towards version 1.0.0!
Firstly, we have dropped support for Python 3.8 and Python 3.9. Python 3.9 is almost at the end of its life cycle, so it's a good time to update our minimum version.
The broker has become a POSITIONAL-ONLY argument. This means that FastStream(broker=broker) is no longer valid. You should always pass the broker as a separate positional argument, like FastStream(brokers), to ensure proper usage.
This is a preparatory step for FastStream(*brokers) support, which will be introduced in 1.0.0.
In 0.6, you can't directly pass custom AsyncAPI options to the FastStream constructor anymore.
app = FastStream( # doesn't work anymore
...,
title="My App",
version="1.0.0",
description="Some description",
)
You need to create a specification object and pass it manually to the constructor.
from faststream import FastStream, AsyncAPI
FastStream(
...
specification=AsyncAPI(
title="My App",
version="1.0.0",
description="Some description",
)
)
Previously, you were able to configure retry attempts for a handler by using the following option:
@broker.subscriber("in", retry=True) # was removed
async def handler(): ...
Unfortunately, this option was a design mistake. We apologize for any confusion it may have caused. Technically, it was just a shortcut to message.nack() on error. We have decided that manual acknowledgement control would be more idiomatic and better for the framework. Therefore, we have provided a new feature in its place: ack_policy control.
@broker.subscriber("test", ack_policy=AckPolicy.ACK_FIRST)
async def handler() -> None: ...
With ack_policy, you can now control the default acknowledge behavior for your handlers. AckPolicy offers the following options:
In addition, we have deprecated a few more options prior to ack_policy.
ack_first=True -> AckPolicy.ACK_FIRSTno_ack=True -> AckPolicy.MANUALWe have made some changes to our Dependency Injection system, so the global context is no longer available.
Currently, you cannot simply import the context from anywhere and use it freely.
from faststeam import context # was removed
Instead, you should create the context in a slightly different way. The FastStream object serves as an entry point for this, so you can place it wherever you need it:
from typing import Annotated
from faststream import Context, ContextRepo, FastStream
from faststream.rabbit import RabbitBroker
broker = RabbitBroker()
app = FastStream(
broker,
context=ContextRepo({
"global_dependency": "value",
}),
)
Everything else about using the context remains the same. You can request it from the context at any place that supports it.
Additionally, Context("broker") and Context("logger") have been moved to the local context. They cannot be accessed from lifespan hooks any longer.
@app.after_startup
async def start(
broker: Broker # does not work anymore
): ...
@router.subscriber
async def handler(
broker: Broker # still working
): ...
This change was also made to support multiple brokers.
Also, we have finalized our Middleware API. It now supports all the features we wanted, and we have no plans to change it anymore. First of all, the BaseMiddleware class constructor requires a context (which is no longer global).
class BaseMiddleware:
def __init__(self, msg: Any | None, context: ContextRepo) -> None:
self.msg = msg
self.context = context
The context is now available as self.context in all middleware methods.
We also changed the publish_scope function signature.
class BaseMiddleware: # old signature
async def publish_scope(
self,
call_next: "AsyncFunc",
msg: Any,
*args: Any,
**kwargs: Any,
) -> Any: ...
Previously, any options passed to brocker.publish("msg", "destination") had to be consumed as *args, **kwargs.
Now, you can consume them all as a single PublishCommand object.
from faststream import PublishCommand
class BaseMiddleware:
async def publish_scope(
self,
call_next: Callable[[PublishCommand], Awaitable[Any]],
cmd: PublishCommand,
) -> Any: ...
Thanks to Python 3.13's TypeVars with defaults, BaseMiddleware becomes a generic class and you can specify the PublishCommand for the broker you want to work with.
from faststream.rabbit import RabbitPublishCommand
class Middleware(BaseMiddleware[RabbitPublishCommand]):
async def publish_scope(
self,
call_next: Callable[[RabbitPublishCommand], Awaitable[Any]],
cmd: RabbitPublishCommand,
) -> Any: ...
Warning: The methods on_consume, after_consume, on_publish and after_publish will be deprecated and removed in version 0.7. Please use consume_scope and publish_scope instead.
In FastStream 0.6 we are using BinaryMessageFormatV1 as a default instead of JSONMessageFormat .
You can find more details in the documentation: https://faststream.ag2.ai/latest/redis/message_format/
AsyncAPI3.0 support – now you can choose between AsyncAPI(schema_version="3.0.0") (default) and AsyncAPI(schema_version="2.6.0") schemas generation
Msgspec native support
from fast_depends.msgspec import MsgSpecSerializer
broker = Broker(serializer=MsgSpecSerializer())
Subscriber iteration support. This features supports all middlewares and other FastStream features.
subscriber = broker.subscriber(..., persistent=False)
await subscriber.start()
async for msg in subscriber:
...
@broker.subscriber(..., filters=...) removedmessage.decoded_body removed, use await message.decode() insteadpublish(..., rpc=True) removed, use broker.request() instead@broker.subscriber(..., reply_config=...) removed, use Response insteadjust by @Samoed in https://github.com/ag2ai/faststream/pull/2436broker.subscriber(persistent=False) argument to control WeakRef behavior by @Lancetnik in https://github.com/ag2ai/faststream/pull/2519Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.48...0.6.0
This is the latest RC version before the stable release. 0.6.0 is scheduled to be released on 10/10/2025.
This is the latest RC version before the stable release. 0.6.0 is scheduled to be released on 10/10/2025.
Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.0rc3...0.6.0rc4
ci: create update release PRs to main: by @Lancetnik in https://github.com/ag2ai/faststream/pull/2490
broker.subscriber(persistent=False) argument to control WeakRef behavior by @Lancetnik in https://github.com/ag2ai/faststream/pull/2519Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.0rc2...0.6.0rc3
fix(aiokafka): AttributeError on first _LoggingListener.on_partitions_assigned by @legau in https://github.com/ag2ai/faststream/pull/2453
Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.0rc1...0.6.0rc2
ci: correct just-install job by @Lancetnik in https://github.com/ag2ai/faststream/pull/2393
just by @Samoed in https://github.com/ag2ai/faststream/pull/2436Full Changelog: https://github.com/ag2ai/faststream/compare/0.6.0rc0...0.6.0rc1
We tried our best to minimize breaking changes, but unfortunately, some aspects were simply not working well. Therefore, we decided to break them in o…
FastStream 0.6 is a significant technical release that aimed to address many of the current project design issues and unlock further improvements on the path to version 1.0.0. We tried our best to minimize breaking changes, but unfortunately, some aspects were simply not working well. Therefore, we decided to break them in order to move forward.
This release includes:
The primary goal of this release is to unlock the path towards further features. Therefore, we are pleased to announce that after this release, we plan to work on MQTT #956 and SQS #794 support and move towards version 1.0.0!
Firstly, we have dropped support for Python 3.8 and Python 3.9. Python 3.9 is almost at the end of its life cycle, so it's a good time to update our minimum version.
The broker has become a POSITIONAL-ONLY argument. This means that FastStream(broker=broker) is no longer valid. You should always pass the broker as a separate positional argument, like FastStream(brokers), to ensure proper usage.
This is a preparatory step for FastStream(*brokers) support, which will be introduced in 1.0.0.
In 0.6, you can't directly pass custom AsyncAPI options to the FastStream constructor anymore.
app = FastStream( # doesn't work anymore
...,
title="My App",
version="1.0.0",
description="Some desctiption",
)
You need to create a specification object and pass it manually to the constructor.
from faststream import FastStream, AsyncAPI
FastStream(
...
specification=AsyncAPI(
title="My App",
version="1.0.0",
description="Some desctiption",
)
)
Previously, you were able to configure retry attempts for a handler by using the following option:
@broker.subscriber("in", retry=True) # was removed
async def handler(): ...
Unfortunately, this option was a design mistake. We apologize for any confusion it may have caused. Technically, it was just a shortcut to message.nack() on error. We have decided that manual acknowledgement control would be more idiomatic and better for the framework. Therefore, we have provided a new feature in its place: ack_policy control.
@broker.subscriber("test", ack_policy=AckPolicy.ACK_FIRST)
async def handler() -> None: ...
With ack_policy, you can now control the default acknowledge behavior for your handlers. AckPolicy offers the following options:
In addition, we have deprecated a few more options prior to ack_policy.
ack_first=True -> AckPolicy.ACK_FIRSTno_ack=True -> AckPolicy.MANUALWe have made some changes to our Dependency Injection system, so the global context is no longer available.
Currently, you cannot simply import the context from anywhere and use it freely.
from faststeam import context # was removed
Instead, you should create the context in a slightly different way. The FastStream object serves as an entry point for this, so you can place it wherever you need it:
from typing import Annotated
from faststream import Context, ContextRepo, FastStream
from faststream.rabbit import RabbitBroker
broker = RabbitBroker()
app = FastStream(
broker,
context=ContextRepo({
"global_dependency": "value",
}),
)
Everything else about using the context remains the same. You can request it from the context at any place that supports it.
Additionally, Context("broker") and Context("logger") have been moved to the local context. They cannot be accessed from lifespan hooks any longer.
@app.after_startup
async def start(
broker: Broker # does not work anymore
): ...
@router.subscriber
async def handler(
broker: Broker # still working
): ...
This change was also made to support multiple brokers.
Also, we have finalized our Middleware API. It now supports all the features we wanted, and we have no plans to change it anymore. First of all, the BaseMiddleware class constructor requires a context (which is no longer global).
class BaseMiddleware:
def __init__(self, msg: Any | None, context: ContextRepo) -> None:
self.msg = msg
self.context = context
The context is now available as self.context in all middleware methods.
We also changed the publish_scope function signature.
class BaseMiddleware: # old signature
async def publish_scope(
self,
call_next: "AsyncFunc",
msg: Any,
*args: Any,
**kwargs: Any,
) -> Any: ...
Previously, any options passed to brocker.publish("msg", "destination") had to be consumed as *args, **kwargs.
Now, you can consume them all as a single PublishCommand object.
from faststream import PublishCommand
class BaseMiddleware:
async def publish_scope(
self,
call_next: Callable[[PublishCommand], Awaitable[Any]],
cmd: PublishCommand,
) -> Any: ...
Thanks to Python 3.13's TypeVars with defaults, BaseMiddleware becomes a generic class and you can specify the PublishCommand for the broker you want to work with.
from faststream.rabbit import RabbitPublishCommand
class Middleware(BaseMiddleware[RabbitPublishCommand]):
async def publish_scope(
self,
call_next: Callable[[RabbitPublishCommand], Awaitable[Any]],
cmd: RabbitPublishCommand,
) -> Any: ...
Warning: The methods on_consume, after_consume, on_publish and after_publish will be deprecated and removed in version 0.7. Please use consume_scope and publish_scope instead.
In FastStream 0.6 we are using BinaryMessageFormatV1 as a default instead of JSONMessageFormat .
You can find more details in the documentation: https://faststream.ag2.ai/latest/redis/message_format/
AsyncAPI3.0 support – now you can choose between AsyncAPI(schema_version="3.0.0") (default) and AsyncAPI(schema_version="2.6.0") schemas generation
Msgspec native support
from fast_depends.msgspec import MsgSpecSerializer
broker = Broker(serializer=MsgSpecSerializer())
Subscriber iteration support. This features supports all middlewares and other FastStream features.
subscriber = broker.subscriber(...)
await subscriber.start()
async for msg in subscriber:
...
@broker.subscriber(..., filters=...) removedmessage.decoded_body removed, use await message.decode() insteadpublish(..., rpc=True) removed, use broker.request() instead@broker.subscriber(..., reply_config=...) removed, use Response insteadFull Changelog: https://github.com/ag2ai/faststream/compare/0.5.48...0.6.0rc0
This release is part of the migration to FastStream 0.6.
This release is part of the migration to FastStream 0.6.
In order to provide great features such as observability and more, FastStream requires the inclusion of additional data in your messages. Redis, on the other hand, allows for the sending of any type of data within a message. Therefore, with this release, we introduce FastStream's own binary message format, which supports any data type you wish to use and can include additional information.
For more information on the message format, please see the documentation
By default, we are still using the JSON message format, but as of version 0.6, the default will change to the binary format. Therefore, you can prepare your services for this change by manually setting a new protocol.
For whole broker:
from faststream.redis import RedisBroker, BinaryMessageFormatV1
# JSONMessageFormat using by default, but it will be deprecated in future updates
broker = RedisBroker(message_format=BinaryMessageFormatV1)
Or for a specifica subscriber / publisher
from faststream.redis import RedisBroker, BinaryMessageFormatV1
broker = RedisBroker()
@broker.subscriber(..., message_format=BinaryMessageFormatV1)
@broker.publisher(..., message_format=BinaryMessageFormatV1)
async def handler(msg):
return msg
Special thanks for @ilya-4real for this great feature!
FastStream require a li
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.47...0.5.48
fix: correct NATS pattern AsyncAPI render by @Lancetnik in https://github.com/ag2ai/faststream/pull/2354
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.46...0.5.47
fix overlapping patterns by @xgemx in https://github.com/ag2ai/faststream/pull/2349
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.45...0.5.46
ci: polish workflows by @Lancetnik in https://github.com/ag2ai/faststream/pull/2331
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.44...0.5.45
introduce stop method on broker and subscriber and deprecate close by @mahenzon in https://github.com/ag2ai/faststream/pull/2328
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.43...0.5.44
fix: deprecate log fmt by @Maclovi in https://github.com/ag2ai/faststream/pull/2240
' in dependency groupds by @sobolevn in https://github.com/ag2ai/faststream/pull/2242Multiple Subscriptions section by @sobolevn in https://github.com/ag2ai/faststream/pull/2250Context Fields Declaration by @sobolevn in https://github.com/ag2ai/faststream/pull/2269Context docs by @sobolevn in https://github.com/ag2ai/faststream/pull/2265publishing/decorator.md by @sobolevn in https://github.com/ag2ai/faststream/pull/2260publishing/index.md docs by @sobolevn in https://github.com/ag2ai/faststream/pull/2259Response class by @anywindblows in https://github.com/ag2ai/faststream/pull/2238serialization/index.md wording and syntax by @sobolevn in https://github.com/ag2ai/faststream/pull/2277Lifespan: Hooks page by @sobolevn in https://github.com/ag2ai/faststream/pull/2284Lifespan: testing by @sobolevn in https://github.com/ag2ai/faststream/pull/2298Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.42...0.5.43
Feature: add deprecate on retry arg by @Flosckow in https://github.com/ag2ai/faststream/pull/2224
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.41...0.5.42
feature: injection FastAPI in StreamMessage scope by @IvanKirpichnikov in https://github.com/ag2ai/faststream/pull/2205
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.40...0.5.41
chore: deprecate connect options by @Lancetnik in https://github.com/ag2ai/faststream/pull/2202
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.39...0.5.40
feat: type client args by @pepellsd in https://github.com/ag2ai/faststream/pull/2165
Full Changelog: https://github.com/ag2ai/faststream/compare/0.5.38...v0.5.39
fix: correct Confluent options name mapping by @Lancetnik in https://github.com/airtai/faststream/pull/2137
Full Changelog: https://github.com/airtai/faststream/compare/0.5.36...0.5.37
fix #2088: respect parsed sasl_mechanism by @Lancetnik in https://github.com/airtai/faststream/pull/2092
CriticalLogMiddleware respect broker log level (#2130) by @fadedDexofan in https://github.com/airtai/faststream/pull/2131Full Changelog: https://github.com/airtai/faststream/compare/0.5.35...0.5.36
Add concurrent-between-partitions kafka subscriber by @Arseniy-Popov in https://github.com/airtai/faststream/pull/2017
Full Changelog: https://github.com/airtai/faststream/compare/0.5.34...0.5.35
fix: when / present in virtual host name and passing as uri by @pepellsd in https://github.com/airtai/faststream/pull/1979
/health by @herotomg in https://github.com/airtai/faststream/pull/2023Full Changelog: https://github.com/airtai/faststream/compare/0.5.33...0.5.34
Just a Confluent & Kafka hotfix. Messages without body (with key only) parsing correctly now.
Just a Confluent & Kafka hotfix. Messages without body (with key only) parsing correctly now.
Full Changelog: https://github.com/airtai/faststream/compare/0.5.32...0.5.33
Thanks to @Flosckow one more time for a new release! Now you have an ability to consume Confluent messages (in autocommit mode) concurrently!
Thanks to @Flosckow one more time for a new release! Now you have an ability to consume Confluent messages (in autocommit mode) concurrently!
from faststream.confluent import KafkaBroker
broker = KafkaBroker()
@broker.subscriber("topic", max_workers=10)
async def handler():
"""Using `max_workers` option you can process up to 10 messages by one subscriber concurrently"""
Also, thanks to @Sehat1137 for his ASGI CLI support bugfixes
Full Changelog: https://github.com/airtai/faststream/compare/0.5.31...0.5.32
Well, you (community) made a new breathtaken release for us! Thanks to all of this release contributors.
Well, you (community) made a new breathtaken release for us! Thanks to all of this release contributors.
Special thanks to @Flosckow . He promotes a new perfect feature - concurrent Kafka subscriber (with autocommit mode)
from faststream.kafka import KafkaBroker
broker = KafkaBroker()
@broker.subscriber("topic", max_workers=10)
async def handler():
"""Using `max_workers` option you can process up to 10 messages by one subscriber concurrently"""
Also, thanks to @Sehat1137 with his ASGI CLI start fixins - now you can use FastStream CLI to scale your AsgiFastStream application by workers
faststream run main:asgi --workers 2
There are a lot of other incredible changes you made:
Full Changelog: https://github.com/airtai/faststream/compare/0.5.30...0.5.31
Introducing FastStream Guru on Gurubase.io by @kursataktas in https://github.com/airtai/faststream/pull/1903
nkeys_seed_str as argument for NATS broker. by @Drakorgaur in https://github.com/airtai/faststream/pull/1908Full Changelog: https://github.com/airtai/faststream/compare/0.5.29...0.5.30
feat: add explicit message source enum by @Lancetnik in https://github.com/airtai/faststream/pull/1866
fake_context if not needed by @sobolevn in https://github.com/airtai/faststream/pull/1877Full Changelog: https://github.com/airtai/faststream/compare/0.5.28...0.5.29
There were a lot of time since 0.5.7 OpenTelemetry release and now we completed Observability features we planned! FastStream supports Prometheus metr
There were a lot of time since 0.5.7 OpenTelemetry release and now we completed Observability features we planned! FastStream supports Prometheus metrics in a native way!
Special thanks to @roma-frolov and @draincoder (again) for it!
To collect Prometheus metrics for your FastStream application you just need to install special distribution
pip install faststream[prometheus]
And use PrometheusMiddleware. Also, it could be helpful to use our ASGI to serve metrics endpoint in the same app.
from prometheus_client import CollectorRegistry, make_asgi_app
from faststream.asgi import AsgiFastStream
from faststream.nats import NatsBroker
from faststream.nats.prometheus import NatsPrometheusMiddleware
registry = CollectorRegistry()
broker = NatsBroker(
middlewares=(
NatsPrometheusMiddleware(registry=registry),
)
)
app = AsgiFastStream(
broker,
asgi_routes=[
("/metrics", make_asgi_app(registry)),
]
)
Moreover, we have a ready-to-use Grafana dashboard you can just import and use!
To find more information about Prometheus support, just visit our documentation.
Full Changelog: https://github.com/airtai/faststream/compare/0.5.27...0.5.28
fix: anyio major version parser by @dotX12 in https://github.com/airtai/faststream/pull/1850
Full Changelog: https://github.com/airtai/faststream/compare/0.5.26...0.5.27
This it the official Python 3.13 support! Now, FastStream works (and tested) at Python 3.8 - 3.13 versions!
This it the official Python 3.13 support! Now, FastStream works (and tested) at Python 3.8 - 3.13 versions!
Warning: Python3.8 is EOF since 3.13 release and we plan to drop it support in FastStream 0.6.0 version.
Also, current release has little bugfixes related to CLI and AsyncAPI schema.
Full Changelog: https://github.com/airtai/faststream/compare/0.5.25...0.5.26
fix: CLI hotfix by @Lancetnik in https://github.com/airtai/faststream/pull/1816
Full Changelog: https://github.com/airtai/faststream/compare/0.5.24...0.5.25
Replace while-sleep with Event by @Olegt0rr in https://github.com/airtai/faststream/pull/1683
Full Changelog: https://github.com/airtai/faststream/compare/0.5.23...0.5.24
We made last release just a few days ago, but there are some big changes here already!
We made last release just a few days ago, but there are some big changes here already!
First of all - you can't use faststream run ... command without pip install faststream[cli] distribution anymore. It was made to minify default (and production) distribution by removing typer (rich and click) dependencies. CLI is a development-time feature, so if you don't need - just don't install! Special thanks to @RubenRibGarcia for this change
The next big change - Kafka publish confirmations by default! Previous FastStream version was working in publish & forgot style, but the new one blocks your broker.publish(...) call until Kafka confirmation frame received. To fallback to previous logic just use a new flag broker.publish(..., no_confirm=True)
Also, we made one more step forward to our 1.0.0 features plan! @KrySeyt implements get_one feature. Now you can use any broker subscriber to get messages in imperative style:
subscriber = broker.subscriber("in")
...
msg = await subscriber.get_one(timeout=5.0)
Big thanks to all new and old contributors who makes such a great release!
broker.subscriber().get_one() by @KrySeyt in https://github.com/airtai/faststream/pull/1726Full Changelog: https://github.com/airtai/faststream/compare/0.5.22...0.5.23
fix: FastAPI 0.112.4+ compatibility by @Lancetnik in https://github.com/airtai/faststream/pull/1766
Full Changelog: https://github.com/airtai/faststream/compare/0.5.21...0.5.22
feat (#1168): allow include regular router to FastAPI integration by @Lancetnik in https://github.com/airtai/faststream/pull/1747
Full Changelog: https://github.com/airtai/faststream/compare/0.5.20...0.5.21
Refactor: change publisher fake subscriber generation logic by @Lancetnik in https://github.com/airtai/faststream/pull/1729
Full Changelog: https://github.com/airtai/faststream/compare/0.5.19...0.5.20
The current release is planned as a latest feature release before 0.6.0. All other 0.5.19+ releases will contain only minor bugfixes and all the team
The current release is planned as a latest feature release before 0.6.0. All other 0.5.19+ releases will contain only minor bugfixes and all the team work will be focused on next major one.
There a lot of changes we want to present you now though!
Our old broker.publish(..., rpc=True) implementation was very limited and ugly. Now we present you a much suitable way to do the same thing - broker.request(...)
from faststream import FastStream
from faststream.nats import NatsBroker, NatsResponse, NatsMessage
broker = NatsBroker()
@broker.subscriber("test")
async def echo_handler(msg):
return NatsResponse(msg, headers={"x-token": "some-token"})
@app.after_startup
async def test():
# The old implementation was returning just a message body,
# so you wasn't be able to check response headers, etc
msg_body: str = await broker.publish("ping", "test", rpc=True)
assert msg_body == "ping"
# Now request return the whole message and you can validate any part of it
# moreover it triggers all your middlewares
response: NatsMessage = await broker.request("ping", "test")
Community asked and community did! Sorry, we've been putting off this job for too long. Thanks for @Rusich90 to help us! Now you can wrap your application by a suitable exception handlers. Just check the new documentation to learn more.
Also, there are a lot of minor changes you can find below. Big thanks to all our old and new contributors! You are amazing ones!
Full Changelog: https://github.com/airtai/faststream/compare/0.5.18...0.5.19
Added additional parameters to HandlerException by @ulbwa in https://github.com/airtai/faststream/pull/1659
Full Changelog: https://github.com/airtai/faststream/compare/0.5.17...0.5.18
Just a hotfix for the following case:
Just a hotfix for the following case:
@broker.subscriber(...)
async def handler():
return NatsResponse(...)
await broker.publish(..., rpc=True)
Full Changelog: https://github.com/airtai/faststream/compare/0.5.16...0.5.17
Well, seems like it is the biggest patch release ever 😃
Well, seems like it is the biggest patch release ever 😃
First of all, thanks to all new contributors, who helps us to improve the project! They made a huge impact to this release by adding new Kafka security mechanisms and extend Response API - now you can use broker.Response to publish detail information from handler
@broker.subscriber("in")
@broker.publisher("out")
async def handler(msg):
return Response(msg, headers={"response_header": "Hi!"}) # or KafkaResponse, etc
Also, we added a new huge feature - ASGI support!
Nope, we are not HTTP-framework now, but it is a little ASGI implementation to provide you with an ability to host documentation, use k8s http-probes and serve metrics in the same with you broker runtime without any dependencies.
You just need to use AsgiFastStream class
from faststream.nats import NatsBroker
from faststream.asgi import AsgiFastStream, make_ping_asgi
from prometheus_client import make_asgi_app
from prometheus_client.registry import CollectorRegistry
broker = NatsBroker()
prometheus_registry = CollectorRegistry()
app = AsgiFastStream(
broker,
asyncapi_path="/docs",
asgi_routes=[
("/health", make_ping_asgi(broker, timeout=5.0)),
("/metrics", make_asgi_app(registry=prometheus_registry))
]
)
And then you can run it like a regular ASGI app
uvicorn main:app
One more thing - manual topic partition assignment for Confluent. We have it already for aiokafka, but missed it here... Now it was fixed!
from faststream.confluent import TopicPartition
@broker.subscriber(partitions=[
TopicPartition("test-topic", partition=0),
])
async def handler():
...
fail_fast option in #1647NatsRouter subjects prefixes behaviorFull Changelog: https://github.com/airtai/faststream/compare/0.5.15...0.5.16
Finally, FastStream has a Kafka pattern subscription! This is another step forward in our Roadmap moving us to 0.6.0 and futher!
Finally, FastStream has a Kafka pattern subscription! This is another step forward in our Roadmap moving us to 0.6.0 and futher!
from faststream import Path
from faststream.kafka import KafkaBroker
broker = KafkaBroker()
@broker.subscriber(pattern="logs.{level}")
async def base_handler(
body: str,
level: str = Path(),
):
...
Also, all brokers now supports a new ping method to check real broker connection
is_connected: bool = await broker.ping()
This is a little, but important change for K8S probes support
More other there are a lot of bugfixes and improvements from our contributors! Thanks to all of these amazing people!
Full Changelog: https://github.com/airtai/faststream/compare/0.5.14...0.5.15
Update Release Notes for 0.5.13 by @faststream-release-notes-updater in https://github.com/airtai/faststream/pull/1548
Full Changelog: https://github.com/airtai/faststream/compare/0.5.13...0.5.14
feat: nats filter JS subscription support by @Lancetnik in https://github.com/airtai/faststream/pull/1519
Full Changelog: https://github.com/airtai/faststream/compare/0.5.12...0.5.13
Now, FastStream provides users with the ability to pass the config dictionary to confluent-kafka-python for greater customizability. The following exa
Now, FastStream provides users with the ability to pass the config dictionary to confluent-kafka-python for greater customizability. The following example sets the parameter topic.metadata.refresh.fast.interval.ms's value to 300 instead of the default value 100 via the config parameter.
from faststream import FastStream
from faststream.confluent import KafkaBroker
config = {"topic.metadata.refresh.fast.interval.ms": 300}
broker = KafkaBroker("localhost:9092", config=config)
app = FastStream(broker)
Full Changelog: https://github.com/airtai/faststream/compare/0.5.11...0.5.12
Update Release Notes for 0.5.10 by @faststream-release-notes-updater in https://github.com/airtai/faststream/pull/1482
Full Changelog: https://github.com/airtai/faststream/compare/0.5.10...0.5.11
Your coding agent can read these notes before it upgrades. Set up the MCP server →