confluent-kafka-dotnet
Show / Hide Table of Contents

Class AsyncSerde<TParsedSchema>

Inheritance
object
AsyncSerde<TParsedSchema>
AsyncDeserializer<T, TParsedSchema>
AsyncSerializer<T, TParsedSchema>
Implements
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
Assembly: Confluent.SchemaRegistry.dll
Syntax
public abstract class AsyncSerde<TParsedSchema> : IClusterIdAware, ISerdeDisposable
Type Parameters
Name Description
TParsedSchema

Constructors

AsyncSerde(ISchemaRegistryClient, SerdeConfig, RuleRegistry)

Declaration
protected AsyncSerde(ISchemaRegistryClient schemaRegistryClient, SerdeConfig config, RuleRegistry ruleRegistry = null)
Parameters
Type Name Description
ISchemaRegistryClient schemaRegistryClient
SerdeConfig config
RuleRegistry ruleRegistry

Fields

associatedNameStrategy

EXPERIMENTAL: subject to change or removal.

The associated subject name strategy backing subjectNameStrategy, when that strategy is Associated. Retained so that the Kafka cluster id resolver can be supplied after construction.

Declaration
protected AssociatedNameStrategy associatedNameStrategy
Field Value
Type Description
AssociatedNameStrategy

latestCompatibilityStrict

Declaration
protected bool latestCompatibilityStrict
Field Value
Type Description
bool

ownsSchemaRegistryClient

EXPERIMENTAL: subject to change or removal.

Whether schemaRegistryClient was constructed by a serde builder rather than supplied by the application, and is therefore disposed along with this instance.

Declaration
protected bool ownsSchemaRegistryClient
Field Value
Type Description
bool

ruleRegistry

Declaration
protected RuleRegistry ruleRegistry
Field Value
Type Description
RuleRegistry

schemaRegistryClient

Declaration
protected ISchemaRegistryClient schemaRegistryClient
Field Value
Type Description
ISchemaRegistryClient

serdeMutex

Declaration
protected SemaphoreSlim serdeMutex
Field Value
Type Description
SemaphoreSlim

subjectNameStrategy

Declaration
protected AsyncSubjectNameStrategyDelegate subjectNameStrategy
Field Value
Type Description
AsyncSubjectNameStrategyDelegate

useLatestVersion

Declaration
protected bool useLatestVersion
Field Value
Type Description
bool

useLatestWithMetadata

Declaration
protected IDictionary<string, string> useLatestWithMetadata
Field Value
Type Description
IDictionary<string, string>

useSchemaId

Declaration
protected int useSchemaId
Field Value
Type Description
int

validationRulesExecution

Declaration
protected ValidationRulesExecution validationRulesExecution
Field Value
Type Description
ValidationRulesExecution

validationRulesFailFast

Declaration
protected bool validationRulesFailFast
Field Value
Type Description
bool

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 virtual void DisposeOwnedResources()

ExecuteMigrations(IList<Migration>, bool, string, string, Headers, object)

Declaration
protected Task<object> ExecuteMigrations(IList<Migration> migrations, bool isKey, string subject, string topic, Headers headers, object message)
Parameters
Type Name Description
IList<Migration> migrations
bool isKey
string subject
string topic
Headers headers
object message
Returns
Type Description
Task<object>

ExecuteRules(bool, string, string, Headers, RuleMode, Schema, Schema, object, FieldTransformer)

Execute rules

Declaration
protected Task<object> ExecuteRules(bool isKey, string subject, string topic, Headers headers, RuleMode ruleMode, Schema source, Schema target, object message, FieldTransformer fieldTransformer)
Parameters
Type Name Description
bool isKey
string subject
string topic
Headers headers
RuleMode ruleMode
Schema source
Schema target
object message
FieldTransformer fieldTransformer
Returns
Type Description
Task<object>
Exceptions
Type Condition
RuleConditionException
ArgumentException

ExecuteRules(bool, string, string, Headers, RulePhase, RuleMode, Schema, Schema, object, FieldTransformer)

Execute rules

Declaration
protected Task<object> ExecuteRules(bool isKey, string subject, string topic, Headers headers, RulePhase rulePhase, RuleMode ruleMode, Schema source, Schema target, object message, FieldTransformer fieldTransformer)
Parameters
Type Name Description
bool isKey
string subject
string topic
Headers headers
RulePhase rulePhase
RuleMode ruleMode
Schema source
Schema target
object message
FieldTransformer fieldTransformer
Returns
Type Description
Task<object>
Exceptions
Type Condition
RuleConditionException
ArgumentException

GetMigrations(string, Schema, RegisteredSchema)

Declaration
protected Task<IList<Migration>> GetMigrations(string subject, Schema writer, RegisteredSchema readerSchema)
Parameters
Type Name Description
string subject
Schema writer
RegisteredSchema readerSchema
Returns
Type Description
Task<IList<Migration>>

GetParsedSchema(Schema)

Declaration
protected Task<TParsedSchema> GetParsedSchema(Schema schema)
Parameters
Type Name Description
Schema schema
Returns
Type Description
Task<TParsedSchema>

GetReaderSchema(string, Schema)

Declaration
protected Task<RegisteredSchema> GetReaderSchema(string subject, Schema schema = null)
Parameters
Type Name Description
string subject
Schema schema
Returns
Type Description
Task<RegisteredSchema>

GetSubjectName(string, bool, string)

Declaration
protected Task<string> GetSubjectName(string topic, bool isKey, string recordType)
Parameters
Type Name Description
string topic
bool isKey
string recordType
Returns
Type Description
Task<string>

GetValidationExecutor()

Returns the executor used to evaluate inline validation rules, or throws when validation is enabled but no rules package has registered one.

Declaration
protected IValidationRuleExecutor GetValidationExecutor()
Returns
Type Description
IValidationRuleExecutor

GetWriterSchema(string, SchemaId, string)

Declaration
protected Task<(Schema, TParsedSchema)> GetWriterSchema(string subject, SchemaId writerId, string format = null)
Parameters
Type Name Description
string subject
SchemaId writerId
string format
Returns
Type Description
Task<(Schema, TParsedSchema)>

IgnoreReference(string)

Declaration
protected virtual bool IgnoreReference(string name)
Parameters
Type Name Description
string name
Returns
Type Description
bool

ParseSchema(Schema)

Declaration
protected abstract Task<TParsedSchema> ParseSchema(Schema schema)
Parameters
Type Name Description
Schema schema
Returns
Type Description
Task<TParsedSchema>

ResolveReferences(Schema)

Declaration
protected Task<IDictionary<string, string>> ResolveReferences(Schema schema)
Parameters
Type Name Description
Schema schema
Returns
Type Description
Task<IDictionary<string, string>>

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.

ValidationEnabled(ValidationRulesExecution?)

Returns true when inline validation rules should run at the given phase.

Pass no phase when there is a single validation point — a serialization path that applies no domain rules has nothing to run before or after, so any enabled mode validates there.

Declaration
protected bool ValidationEnabled(ValidationRulesExecution? phase = null)
Parameters
Type Name Description
ValidationRulesExecution? phase
Returns
Type Description
bool

Implements

IClusterIdAware
ISerdeDisposable
In this article