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[].
Inherited Members
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
|
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. |