confluent-kafka-dotnet
Show / Hide Table of Contents

Class AvroSerializer<T>

(async) Avro serializer. Use this serializer with GenericRecord, types generated using the avrogen.exe tool or one of the following primitive types: int, long, float, double, boolean, string, byte[].

Inheritance
object
AvroSerializer<T>
Implements
IAsyncSerializer<T>
IClusterIdAware
ISerdeDisposable
Inherited Members
object.Equals(object)
object.Equals(object, object)
object.GetHashCode()
object.GetType()
object.MemberwiseClone()
object.ReferenceEquals(object, object)
object.ToString()
Namespace: Confluent.SchemaRegistry.Serdes
Assembly: Confluent.SchemaRegistry.Serdes.Avro.dll
Syntax
public class AvroSerializer<T> : IAsyncSerializer<T>, IClusterIdAware, ISerdeDisposable
Type Parameters
Name Description
T
Remarks

Serialization format: byte 0: Magic byte use to identify the protocol format. bytes 1-4: Unique global id of the Avro schema that was used for encoding (as registered in Confluent Schema Registry), big endian. following bytes: The serialized data.

Constructors

AvroSerializer(ISchemaRegistryClient, AvroSerializerConfig, RuleRegistry)

Initialize a new instance of the AvroSerializer class. When passed as a parameter to the Confluent.Kafka.Producer constructor, the following configuration properties will be extracted from the producer's configuration property collection:

avro.serializer.buffer.bytes (default: 128) - Initial size (in bytes) of the buffer used for message serialization. Use a value high enough to avoid resizing the buffer, but small enough to avoid excessive memory use. Inspect the size of the byte array returned by the Serialize method to estimate an appropriate value. Note: each call to serialize creates a new buffer.

avro.serializer.auto.register.schemas (default: true) - true if the serializer should attempt to auto-register unrecognized schemas with Confluent Schema Registry, false if not.

Declaration
public AvroSerializer(ISchemaRegistryClient schemaRegistryClient, AvroSerializerConfig config = null, RuleRegistry ruleRegistry = null)
Parameters
Type Name Description
ISchemaRegistryClient schemaRegistryClient

An implementation of ISchemaRegistryClient used for communication with Confluent Schema Registry.

AvroSerializerConfig config

Serializer configuration properties (refer to AvroSerializerConfig)

RuleRegistry ruleRegistry

AvroSerializer(ISchemaRegistryClient, IEnumerable<KeyValuePair<string, string>>)

Initialize a new instance of the AvroSerializer class.

Declaration
[Obsolete("Superseded by AvroSerializer(ISchemaRegistryClient, AvroSerializerConfig)")]
public AvroSerializer(ISchemaRegistryClient schemaRegistryClient, IEnumerable<KeyValuePair<string, string>> config)
Parameters
Type Name Description
ISchemaRegistryClient schemaRegistryClient
IEnumerable<KeyValuePair<string, string>> config

Fields

DefaultInitialBufferSize

The default initial size (in bytes) of buffers used for message serialization.

Declaration
public const int DefaultInitialBufferSize = 1024
Field Value
Type Description
int

Methods

DisposeOwnedResources()

Release the resources this instance created itself.

Resources supplied by the application are not released. Implementations must tolerate being called more than once.

Declaration
public void DisposeOwnedResources()

SerializeAsync(T, SerializationContext)

Serialize an instance of type T to a byte array in Avro format. The serialized data is preceded by a "magic byte" (1 byte) and the id of the schema as registered in Confluent's Schema Registry (4 bytes, network byte order). This call may block or throw on first use for a particular topic during schema registration.

Declaration
public Task<byte[]> SerializeAsync(T value, SerializationContext context)
Parameters
Type Name Description
T value

The value to serialize.

SerializationContext context

Context relevant to the serialize operation.

Returns
Type Description
Task<byte[]>

A Task that completes with value serialized as a byte array.

SetClusterIdResolver(Func<Task<string>>)

Supply a resolver for the id of the Kafka cluster the client is connected to.

The resolver returns a task that completes once the client has reached a broker, with null if it cannot do so within the client's timeout; invoking it never blocks the caller. Concurrent invocations share a single resolution. Implementations are expected to invoke it lazily, only when the id is actually needed.

The resolver is bound to the client that supplied it, and throws ObjectDisposedException once that client has been disposed. A serializer or deserializer handed to a producer or consumer must therefore not be used after that client is disposed, unless the cluster id it needs was specified via configuration.

A serializer or deserializer instance holds a single resolver. If one instance is shared by several clients, the resolver supplied last replaces the earlier ones, so every client sharing the instance resolves the cluster id of the last client it was handed to. Do not share an instance between clients connected to different clusters; let the client creates the serde instance through the builder.

Implementations must ignore the resolver when the cluster id is not relevant to their configuration, or when it was specified explicitly via configuration, so that a configured cluster id is never overwritten.

Declaration
public void SetClusterIdResolver(Func<Task<string>> clusterIdResolver)
Parameters
Type Name Description
Func<Task<string>> clusterIdResolver

Resolves the Kafka cluster id.

Implements

IAsyncSerializer<T>
IClusterIdAware
ISerdeDisposable
In this article