Class AsyncSerde<TParsedSchema>
Inheritance
AsyncSerde<TParsedSchema>
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
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
latestCompatibilityStrict
Declaration
protected bool latestCompatibilityStrict
Field Value
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
ruleRegistry
Declaration
protected RuleRegistry ruleRegistry
Field Value
schemaRegistryClient
Declaration
protected ISchemaRegistryClient schemaRegistryClient
Field Value
serdeMutex
Declaration
protected SemaphoreSlim serdeMutex
Field Value
subjectNameStrategy
Declaration
protected AsyncSubjectNameStrategyDelegate subjectNameStrategy
Field Value
useLatestVersion
Declaration
protected bool useLatestVersion
Field Value
Declaration
protected IDictionary<string, string> useLatestWithMetadata
Field Value
useSchemaId
Declaration
protected int useSchemaId
Field Value
validationRulesExecution
Declaration
protected ValidationRulesExecution validationRulesExecution
Field Value
validationRulesFailFast
Declaration
protected bool validationRulesFailFast
Field Value
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()
Declaration
protected Task<object> ExecuteMigrations(IList<Migration> migrations, bool isKey, string subject, string topic, Headers headers, object message)
Parameters
Returns
Declaration
protected Task<object> ExecuteRules(bool isKey, string subject, string topic, Headers headers, RuleMode ruleMode, Schema source, Schema target, object message, FieldTransformer fieldTransformer)
Parameters
Returns
Exceptions
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
Returns
Exceptions
GetMigrations(string, Schema, RegisteredSchema)
Declaration
protected Task<IList<Migration>> GetMigrations(string subject, Schema writer, RegisteredSchema readerSchema)
Parameters
Returns
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
Returns
GetSubjectName(string, bool, string)
Declaration
protected Task<string> GetSubjectName(string topic, bool isKey, string recordType)
Parameters
Returns
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
GetWriterSchema(string, SchemaId, string)
Declaration
protected Task<(Schema, TParsedSchema)> GetWriterSchema(string subject, SchemaId writerId, string format = null)
Parameters
Returns
IgnoreReference(string)
Declaration
protected virtual bool IgnoreReference(string name)
Parameters
| Type |
Name |
Description |
| string |
name |
|
Returns
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
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
Returns
Implements