pub trait Deserializer<T>:
Send
+ Sync
+ 'static {
// Required method
fn deserialize(&self, topic: &str, data: &[u8]) -> Result<T, Error>;
// Provided methods
fn deserialize_headers(
&self,
topic: &str,
_headers: &dyn Headers,
data: &[u8],
) -> Result<T, Error> { ... }
fn deserialize_from_shared(
&self,
topic: &str,
_source: &Bytes,
data: &[u8],
) -> Result<T, Error> { ... }
fn deserialize_from_shared_headers(
&self,
topic: &str,
headers: &dyn Headers,
_source: &Bytes,
data: &[u8],
) -> Result<T, Error> { ... }
fn configure(&mut self, _configs: &HashMap<String, String>, _is_key: bool) { ... }
fn close(&mut self) { ... }
}Expand description
An interface for converting bytes to objects.
Corresponds to Java’s org.apache.kafka.common.serialization.Deserializer<T>.
§Sync, not async
deserialize is intentionally a synchronous fn (not async) so that the
receive path can call it inline against borrowed slices of the
CompletedFetch buffer without allocating a per-record future. See
consumer-threading.md §27 — every record allocates at most the
user-deserialized T; the trait itself adds zero heap allocation per
call.
External lookups inside deserialize: deserializers that depend on
external state (most notably a schema registry) should pre-populate an
in-memory cache before the consumer starts polling. For rare blocking
calls inside deserialize, callers on the multi-thread runtime can wrap
with tokio::task::block_in_place; this is not free and should not be
the per-record default.
§Bounds
Send + Sync + 'static is required because deserializer instances are
stored as Box<dyn Deserializer<T>> inside Deserializers<K, V>, which
is shared across tasks via Arc<Deserializers<K, V>> (the consumer’s app
side and its background task hold the same Arc). Arc<T>: Send
requires T: Send + Sync; that requirement transits the Box<dyn>
boundary to the trait.
Required Methods§
Sourcefn deserialize(&self, topic: &str, data: &[u8]) -> Result<T, Error>
fn deserialize(&self, topic: &str, data: &[u8]) -> Result<T, Error>
Deserialize a record value from a byte slice.
Corresponds to Java’s T deserialize(String topic, byte[] data).
§Arguments
topic- topic associated with the datadata- serialized bytes; the receive path slices these from the underlyingCompletedFetchbuffer (consumer-threading.md§27).
§Returns
The deserialized typed object, or a Error if deserialization
fails.
Java returns T directly and accepts a null byte[] returning a
null T. In Rust the receive path never passes a null slice — it
passes an empty slice or skips the call entirely — so the trait
surface only models the happy path. Deserialization errors are
reported via the Result.
Provided Methods§
Sourcefn deserialize_headers(
&self,
topic: &str,
_headers: &dyn Headers,
data: &[u8],
) -> Result<T, Error>
fn deserialize_headers( &self, topic: &str, _headers: &dyn Headers, data: &[u8], ) -> Result<T, Error>
Deserialize a record value with access to its headers.
Corresponds to Java’s
default T deserialize(String topic, Headers headers, byte[] data)
(Deserializer.java:84). The default implementation ignores the headers
and delegates to deserialize.
The translated Java overloads intersect on {topic, data}, which is
exactly deserialize(String, byte[]) (:64) — so that one keeps the
plain name and this one is suffixed with the parameter that distinguishes
it (CLAUDE.md §2). Java’s third overload, deserialize(String, Headers, ByteBuffer) (:113), is conversion sugar over this one and has no Rust
counterpart: Rust has no ByteBuffer, and both forms would translate to
the same &[u8] signature.
Override this method in custom deserializer implementations that need to inspect headers during deserialization (for example, schema registry integration).
Takes &dyn Headers (the trait, not the concrete RecordHeaders
struct) to mirror Java’s signature exactly — Java’s parameter type
is the Headers interface, which lets test doubles and alternative
Headers implementations interoperate with custom deserializers.
Deserialize from a byte slice that is a subslice of a refcounted source
buffer, enabling zero-copy T for byte-typed deserializers.
source is the [bytes::Bytes] that owns the whole records buffer for
the current batch (on the receive path); data is the key/value
subslice of source for this record. The default implementation ignores
source and delegates to deserialize — so
existing deserializers are unaffected. BytesDeserializer overrides it
to return source.slice_ref(data), a zero-copy refcounted slice instead
of an owned copy (consumer-threading.md §27).
Header-aware variant of
deserialize_from_shared.
The default implementation delegates to
deserialize_headers,
ignoring source. This preserves the header-inspection behavior of any
deserializer that overrides deserialize_headers (e.g. schema
registry) even when called on the shared-buffer receive path — at the
cost of the copy fallback. Byte-typed deserializers that want zero-copy
override this method directly (see BytesDeserializer).