Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
fe2d1c1
feat(codec): add destination parameter to encode
ce1ebrimbor Jun 3, 2026
398ff19
docs(codec): add gzip codec examples and documentation
ce1ebrimbor Jun 3, 2026
0a13d8a
feat(codec): change encode signature to accept PublishCommand
ce1ebrimbor Jun 4, 2026
f0fa7e4
docs(codec): replace examples with Schema Registry codec
ce1ebrimbor Jun 4, 2026
17935db
feat(codec): pass PublishCommand directly to codec.encode() in Redis …
ce1ebrimbor Aug 6, 2026
29084db
refactor(confluent): move function-level imports to module top
ce1ebrimbor Aug 6, 2026
f0cd4b4
refactor(kafka): move function-level imports to module top
ce1ebrimbor Aug 6, 2026
1a9fd66
refactor(redis,rabbit): move function-level imports to module top
ce1ebrimbor Aug 6, 2026
ab605d4
chore: remove .omo/ from tracking
ce1ebrimbor Aug 6, 2026
5d6dafe
refactor: rename _BaseCmd to _BasePublishCommand for consistency
ce1ebrimbor Aug 6, 2026
03fa2a5
style: fix linting errors
ce1ebrimbor Aug 6, 2026
6ff15ca
refactor(rabbit): remove redundant destination param from encode_message
ce1ebrimbor Aug 6, 2026
454d3a2
fix(rabbit): make cmd required in _publish to match encode_message si…
ce1ebrimbor Aug 6, 2026
a2f7135
refactor(codec): replace tuple return with EncodedMessage dataclass
ce1ebrimbor Aug 7, 2026
ff5b440
style: format codec example
ce1ebrimbor Aug 7, 2026
c99254a
refactor: remove _BasePublishCommand alias, use PublishCommand directly
ce1ebrimbor Aug 8, 2026
81f3b93
fix: clean up suggestions
ce1ebrimbor Aug 8, 2026
8006c97
docs: use mkdocstrings for CodecProto, fix return type description
ce1ebrimbor Aug 8, 2026
9735038
refactor(rabbit): simplify encode_message to extract params from Rabb…
ce1ebrimbor Aug 8, 2026
1cd5447
ci: add _typos.toml to suppress WellFREEzZ false positive
ce1ebrimbor Aug 9, 2026
5165587
fix: use super().encode reference test batch brokers
ce1ebrimbor Aug 9, 2026
08907ba
refactor(tests): use super().encode via encode_item in batch codec tests
ce1ebrimbor Aug 9, 2026
1952e35
refactor: re-export codec types from faststream.message, DefaultCodec…
ce1ebrimbor Aug 14, 2026
c6e3886
refactor: deprecate parser/decoder params with DeprecationWarning
ce1ebrimbor Aug 14, 2026
d565255
refactor: move warn_deprecated_param import to top of usecase.py
ce1ebrimbor Aug 14, 2026
3db8aeb
refactor: deprecate only decoder param, not parser
ce1ebrimbor Aug 14, 2026
04993de
refactor: remove unused warn_deprecated_param helper
ce1ebrimbor Aug 15, 2026
df5e22f
fix: linter
ce1ebrimbor Aug 15, 2026
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
4 changes: 4 additions & 0 deletions _typos.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
[default]
extend-ignore-re = [
"WellFREEzZ",
]
59 changes: 59 additions & 0 deletions docs/docs/en/getting-started/serialization/codec.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
---
# 0.5 - API
# 2 - Release
# 3 - Contributing
# 5 - Template Page
# 10 - Default
search:
boost: 10
---

# Custom Codec

A codec provides a unified interface for both encoding (publishing) and decoding (consuming) messages. Unlike the older `decoder=` approach, a codec handles both directions in a single class.

## Protocol

Implement the `CodecProto` interface to create a custom codec:

::: faststream._internal.parser.CodecProto

- **`decode`** — receives a `StreamMessage` with raw bytes in `msg.body` and returns the decoded Python value.
- **`encode`** — receives a `PublishCommand` containing the message body, destination, and headers. Returns an `EncodedMessage` dataclass with `body: bytes` and `content_type: str | None`. Access the payload via `cmd.body` and the target topic/subject/queue via `cmd.destination`.

If no codec is set, `DefaultCodec` is used automatically. It handles JSON objects, plain text, and raw bytes.

## Example: Schema Registry

A Confluent Avro codec that encodes and decodes messages using the [Confluent wire format](https://docs.confluent.io/platform/current/schema-registry/fundamentals/serdes-develop/index.html#wire-format){target="_blank"} (magic byte + schema ID + Avro payload). Requires `fastavro` and `confluent-kafka`:

```bash
pip install fastavro confluent-kafka
```

```python linenums="1" hl_lines="22-66 68-76"
{!> docs_src/getting_started/serialization/codec_schema_registry_kafka.py !}
```

!!! note
The codec fetches and caches schemas from the registry at startup and on first encounter. The `subject` follows Confluent's naming convention: `{topic}-value`.

## Priority

You can set a codec at the broker level or override it per subscriber. The subscriber-level codec always wins:

```python
broker = KafkaBroker(codec=BrokerCodec())

@broker.subscriber("test", codec=SubscriberCodec()) # ← this wins
async def handle(body: str) -> None:
...

# If no codec is set at any level, DefaultCodec is used (JSON/text/bytes)
```

## Compatibility

- **`codec=` and `parser=`** work together. The parser controls how the raw broker message is parsed into a `StreamMessage`; the codec then decodes or encodes the body.
- **`codec=` and `decoder=`** cannot be used together. Specifying both raises a `ValueError`.
- For the legacy `decoder=` approach, see [Custom Decoder](./decoder.md){.internal-link}.
5 changes: 5 additions & 0 deletions docs/docs/en/getting-started/serialization/decoder.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,11 @@ search:
boost: 10
---

!!! warning "Superseded by Codec"
The `decoder=` parameter has been superseded by the new **Codec** system, which handles both encoding and decoding in a single interface. See [Custom Codec](./codec.md){.internal-link} for the recommended approach.

Note: `codec=` and `decoder=` cannot be used together — specifying both will raise a `ValueError`.

# Custom Decoder

At this stage, the body of a **StreamMessage** is transformed into the format that it will take when it enters your handler function. This stage is the one you will need to redefine more often.
Expand Down
1 change: 1 addition & 0 deletions docs/docs/navigation_template.txt
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ search:
- [Context](getting-started/context.md)
- [Custom Serialization](getting-started/serialization/index.md)
- [Parser](getting-started/serialization/parser.md)
- [Codec](getting-started/serialization/codec.md)
- [Decoder](getting-started/serialization/decoder.md)
- [Examples](getting-started/serialization/examples.md)
- [Lifespan](getting-started/lifespan/index.md)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
import io
import json
import struct
from typing import TYPE_CHECKING, Any, Dict

import fastavro
from confluent_kafka.schema_registry import SchemaRegistryClient

from faststream import FastStream
from faststream.kafka import KafkaBroker

if TYPE_CHECKING:
from fast_depends.library.serializer import SerializerProto

from faststream._internal.basic_types import DecodedMessage
from faststream.message import StreamMessage
from faststream.response.response import PublishCommand

HEADER = struct.Struct(">bI") # magic byte (0) + 4-byte schema ID


class SchemaRegistryCodec:
def __init__(
self,
registry_url: str,
topics: Dict[str, int],
) -> None:
self._client = SchemaRegistryClient({"url": registry_url})
self._schema_cache: Dict[int, Any] = {}
self._topic_schemas: Dict[str, tuple[int, Any]] = {}

for topic, version in topics.items():
subject = f"{topic}-value"
meta = self._client.get_version(subject, version)
schema = fastavro.parse_schema(json.loads(meta.schema.schema_str))
self._topic_schemas[topic] = (meta.schema_id, schema)
self._schema_cache[meta.schema_id] = schema

def _get_schema(self, schema_id: int) -> Any:
if schema_id not in self._schema_cache:
raw = self._client.get_schema(schema_id)
self._schema_cache[schema_id] = fastavro.parse_schema(
json.loads(raw.schema_str)
)
return self._schema_cache[schema_id]

async def decode(self, msg: "StreamMessage[Any]") -> "DecodedMessage":
schema_id = int.from_bytes(msg.body[1:5], byteorder="big")
schema = self._get_schema(schema_id)
decoded: dict[str, Any] = fastavro.schemaless_reader(
io.BytesIO(msg.body[5:]), schema
)
return decoded # type: ignore[return-value]

async def encode(
self,
cmd: "PublishCommand",
serializer: "SerializerProto | None" = None,
) -> "EncodedMessage":
from faststream._internal.parser import EncodedMessage

schema_id, schema = self._topic_schemas[cmd.destination]
body = cmd.body
data = body.model_dump(mode="json") if hasattr(body, "model_dump") else body
buf = io.BytesIO()
buf.write(HEADER.pack(0, schema_id))
fastavro.schemaless_writer(buf, schema, data)
return EncodedMessage(body=buf.getvalue(), content_type="application/avro")


codec = SchemaRegistryCodec(
registry_url="http://localhost:8081",
topics={
"orders": 1,
"users": 2,
},
)
broker = KafkaBroker(codec=codec)
app = FastStream(broker)


@broker.subscriber("orders")
async def handle_order(body: dict[str, Any]) -> None: ...


@broker.subscriber("users")
async def handle_user(body: dict[str, Any]) -> None: ...


@app.after_startup
async def test() -> None:
await broker.publish({"order_id": "123", "amount": 99.99}, "orders")
await broker.publish({"name": "John", "age": 25}, "users")
10 changes: 10 additions & 0 deletions faststream/_internal/configs/broker.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import warnings
from collections.abc import Iterable, Sequence
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any, Generic, Optional, Union
Expand Down Expand Up @@ -40,6 +41,15 @@ class BrokerConfig:
ack_policy: "AckPolicy" = field(default_factory=lambda: EMPTY)
extra_context: dict[str, Any] = field(default_factory=dict)

def __post_init__(self) -> None:
if self.broker_decoder is not None:
warnings.warn(
"`decoder` parameter is deprecated and will be removed in 1.0.0. "
"Use `codec` with a custom `encode()`/`decode()` method instead.",
DeprecationWarning,
stacklevel=3,
)

def __repr__(self) -> str:
return f"{self.__class__.__name__}(id: {id(self)})"

Expand Down
17 changes: 17 additions & 0 deletions faststream/_internal/endpoint/subscriber/usecase.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import warnings
from abc import abstractmethod
from collections.abc import AsyncIterator, Callable, Iterable, Sequence
from contextlib import AbstractContextManager, AsyncExitStack
Expand Down Expand Up @@ -234,6 +235,14 @@ def add_call(
dependencies_: Iterable["Dependant"],
codec_: Optional["CodecProto"] = None,
) -> Self:
if decoder_ is not None:
warnings.warn(
"`decoder` parameter is deprecated and will be removed in 1.0.0. "
"Use `codec` with a custom `encode()`/`decode()` method instead.",
DeprecationWarning,
stacklevel=3,
)

self._call_options = _CallOptions(
parser=parser_,
decoder=decoder_,
Expand Down Expand Up @@ -283,6 +292,14 @@ def __call__(
"HandlerCallWrapper[P_HandlerParams, T_HandlerReturn]",
],
]:
if decoder is not None:
warnings.warn(
"`decoder` parameter is deprecated and will be removed in 1.0.0. "
"Use `codec` with a custom `encode()`/`decode()` method instead.",
DeprecationWarning,
stacklevel=3,
)

total_deps = (*self._call_options.dependencies, *dependencies)
async_filter: AsyncFilter[StreamMessage[MsgType]] = to_async(filter)

Expand Down
29 changes: 20 additions & 9 deletions faststream/_internal/parser.py
Original file line number Diff line number Diff line change
@@ -1,18 +1,28 @@
from abc import abstractmethod
from collections.abc import Sequence
from dataclasses import dataclass
from typing import TYPE_CHECKING, Any, Protocol, TypeVar, runtime_checkable

from faststream.message.utils import decode_message, encode_message

if TYPE_CHECKING:
from fast_depends.library.serializer import SerializerProto

from faststream._internal.basic_types import DecodedMessage, SendableMessage
from faststream._internal.basic_types import DecodedMessage
from faststream.message import StreamMessage
from faststream.response.response import PublishCommand

MsgType = TypeVar("MsgType")


@dataclass(frozen=True)
class EncodedMessage:
"""Result of codec encoding."""

body: bytes
content_type: str | None = None


class ParserProto(Protocol[MsgType]):
"""Protocol for parsing raw messages into StreamMessage."""

Expand Down Expand Up @@ -64,19 +74,19 @@ async def decode(self, msg: "StreamMessage[Any]") -> "DecodedMessage":
@abstractmethod
async def encode(
self,
msg: "SendableMessage",
cmd: "PublishCommand",
Comment thread
ce1ebrimbor marked this conversation as resolved.
serializer: "SerializerProto | None" = None,
) -> tuple[bytes, str | None]: ...
) -> "EncodedMessage": ...


@runtime_checkable
class BatchCodecProto(Protocol):
@abstractmethod
async def encode_batch(
self,
msgs: Sequence["SendableMessage"],
cmd: "PublishCommand",
serializer: "SerializerProto | None" = None,
) -> list[tuple[bytes, str | None]]: ...
) -> list["EncodedMessage"]: ...

@abstractmethod
async def decode_batch(
Expand All @@ -85,13 +95,14 @@ async def decode_batch(
) -> list["DecodedMessage"]: ...


class DefaultCodec:
class DefaultCodec(CodecProto):
async def decode(self, msg: "StreamMessage[Any]") -> "DecodedMessage":
return decode_message(msg)

async def encode(
self,
msg: "SendableMessage",
cmd: "PublishCommand",
serializer: "SerializerProto | None" = None,
) -> tuple[bytes, str | None]:
return encode_message(msg, serializer)
) -> "EncodedMessage":
body, content_type = encode_message(cmd.body, serializer)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's make encode_message return EncodedMessage to be consistent, because decode_message returns DecodedMessage

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I am not sure this is the right abstraction level, if you look at the decode_message function in faststream/message/utils.py(same module) it returns:

def decode_message(message: "StreamMessage[Any]") -> "DecodedMessage":

this method retruns a different kind of DecodedMessage:

DecodedMessage: TypeAlias = JsonDecodable | JsonArray | JsonTable

codecs' EncodedMessage type does not really belong in this layer (imo).
If we want to make it consistent It will require more work

return EncodedMessage(body=body, content_type=content_type)
34 changes: 22 additions & 12 deletions faststream/confluent/publisher/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from faststream.confluent.parser import AsyncConfluentParser
from faststream.confluent.response import KafkaPublishCommand
from faststream.exceptions import FeatureNotSupportedException
from faststream.response.response import PublishCommand

from .state import EmptyProducerState, ProducerState, RealProducer

Expand Down Expand Up @@ -139,10 +140,13 @@ async def publish(
cmd: "KafkaPublishCommand",
) -> "asyncio.Future[Message | None] | Message | None":
"""Publish a message to a topic."""
if cmd.body is None:
message, content_type = None, None
if cmd.body is None and cmd.key is not None:
body = None
content_type = None
else:
message, content_type = await self.codec.encode(cmd.body, self.serializer)
encoded = await self.codec.encode(cmd, self.serializer)
body = encoded.body
content_type = encoded.content_type

headers_to_send = {
"content-type": content_type or "",
Expand All @@ -151,7 +155,7 @@ async def publish(

return await self._producer.producer.send(
topic=cmd.destination,
value=message,
value=body,
key=cmd.key,
partition=cmd.partition,
timestamp_ms=cmd.timestamp_ms,
Expand All @@ -167,26 +171,32 @@ async def publish_batch(self, cmd: "KafkaPublishCommand") -> None:
headers_to_send = cmd.headers_to_publish()

if isinstance(self.codec, BatchCodecProto):
encoded_batch = await self.codec.encode_batch(
cmd.batch_bodies, self.serializer
)
encoded_batch = await self.codec.encode_batch(cmd, self.serializer)
else:
encoded_batch = [
await self.codec.encode(msg, self.serializer) for msg in cmd.batch_bodies
await self.codec.encode(
PublishCommand(
body=msg,
destination=cmd.destination,
_publish_type=cmd.publish_type,
),
self.serializer,
)
for msg in cmd.batch_bodies
]

for message_position, (message, content_type) in enumerate(encoded_batch):
if content_type:
for message_position, encoded in enumerate(encoded_batch):
if encoded.content_type:
final_headers = {
"content-type": content_type,
"content-type": encoded.content_type,
**headers_to_send,
}
else:
final_headers = headers_to_send.copy()

batch.append(
key=cmd.key_for(message_position),
value=message,
value=encoded.body,
timestamp=cmd.timestamp_ms,
headers=[(i, j.encode()) for i, j in final_headers.items()],
)
Expand Down
Loading
Loading