Repository navigation
feat(codec): wire codecs into publishing flow #2902
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
fe2d1c1
398ff19
0a13d8a
f0fa7e4
17935db
29084db
f0cd4b4
1a9fd66
ab605d4
5d6dafe
03fa2a5
6ff15ca
454d3a2
a2f7135
ff5b440
c99254a
81f3b93
8006c97
9735038
1cd5447
5165587
08907ba
1952e35
c6e3886
d565255
3db8aeb
04993de
df5e22f
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,4 @@ | ||
| [default] | ||
| extend-ignore-re = [ | ||
| "WellFREEzZ", | ||
| ] |
| 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}. |
| 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") |
| 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.""" | ||
|
|
||
|
|
@@ -64,19 +74,19 @@ async def decode(self, msg: "StreamMessage[Any]") -> "DecodedMessage": | |
| @abstractmethod | ||
| async def encode( | ||
| self, | ||
| msg: "SendableMessage", | ||
| cmd: "PublishCommand", | ||
| 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( | ||
|
|
@@ -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) | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Let's make
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 def decode_message(message: "StreamMessage[Any]") -> "DecodedMessage":this method retruns a different kind of DecodedMessage: DecodedMessage: TypeAlias = JsonDecodable | JsonArray | JsonTablecodecs' EncodedMessage type does not really belong in this layer (imo). |
||
| return EncodedMessage(body=body, content_type=content_type) | ||
Uh oh!
There was an error while loading. Please reload this page.