Confluent.Kafka.FactoryExtensions
9.0.0
dotnet add package Confluent.Kafka.FactoryExtensions --version 9.0.0
NuGet\Install-Package Confluent.Kafka.FactoryExtensions -Version 9.0.0
<PackageReference Include="Confluent.Kafka.FactoryExtensions" Version="9.0.0" />
<PackageVersion Include="Confluent.Kafka.FactoryExtensions" Version="9.0.0" />
<PackageReference Include="Confluent.Kafka.FactoryExtensions" />
paket add Confluent.Kafka.FactoryExtensions --version 9.0.0
#r "nuget: Confluent.Kafka.FactoryExtensions, 9.0.0"
#:package Confluent.Kafka.FactoryExtensions@9.0.0
#addin nuget:?package=Confluent.Kafka.FactoryExtensions&version=9.0.0
#tool nuget:?package=Confluent.Kafka.FactoryExtensions&version=9.0.0
Confluent Kafka Extension: Client Factory
An extension of Confluent's .NET Client for Apache Kafka<sup>TM</sup>.
Installation
- Package Manager
Install-Package Confluent.Kafka.FactoryExtensions -Version 8.x.x
- .NET CLI
dotnet add package Confluent.Kafka.FactoryExtensions --version 8.x.x
- PackageReference
<PackageReference Include="Confluent.Kafka.FactoryExtensions" Version="8.x.x" />
Features
- Configure multiple Named Kafka Clients using
Microsoft.Extensions.Configuration.IConfiguration. - Register Kafka Clients in a thread safe Concurrent Dictionary.
- Inject IConsumerFactory and IProducerFactory using
Microsoft.Extensions.DependencyInjection.
Usage
Take a look in the examples directory for example usage.
{
"Kafka": {
"Consumers": {
"Constellation": {
"Topic": "EVENTHUB",
"Config": {
"BootstrapServers": "NAMESPACENAME.servicebus.windows.net:9093",
"ClientId": null,
"Debug": null,
"SaslMechanism": "Plain",
"SaslUsername": "$ConnectionString",
"SecurityProtocol": "SaslSsl",
"SocketKeepaliveEnable": true,
"SaslPassword": "{YOUR.EVENTHUBS.CONNECTION.STRING}",
"GroupId": "CONSUMERGROUP",
"AutoOffsetReset": "Latest",
"BrokerVersionFallback": "1.0.0"
}
},
"Qualification": {
"Topic": "EVENTHUB",
"Config": {
"BootstrapServers": "NAMESPACENAME.servicebus.windows.net:9093",
"ClientId": null,
"Debug": null,
"SaslMechanism": "Plain",
"SaslUsername": "$ConnectionString",
"SecurityProtocol": "SaslSsl",
"SocketKeepaliveEnable": true,
"SaslPassword": "{YOUR.EVENTHUBS.CONNECTION.STRING}",
"GroupId": "CONSUMERGROUP",
"AutoOffsetReset": "Latest",
"BrokerVersionFallback": "1.0.0"
}
}
}
}
}
Important: Replace {YOUR.EVENTHUBS.CONNECTION.STRING} with the connection string for your Event Hubs namespace. For instructions on getting the connection string, see Get an Event Hubs connection string. Here's an example configuration: "SaslPassword" : "Endpoint=sb://mynamespace.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=XXXXXXXXXXXXXXXX";
DI Configuration
Add all .json settings files to IConfigurationBuilder,
get configuration section var configuration = hostContext.Configuration.GetSection(nameof(Kafka));
and register the client factories in DI with services.TryAddKafkaFactories(configuration);
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureAppConfiguration((context, builder) =>
{
var configuration = builder.Build();
var configSubPath = configuration.GetValue<string>("CONFIG_SUB_PATH");
var directoryContents = context.HostingEnvironment.ContentRootFileProvider.GetDirectoryContents(configSubPath);
foreach (var file in directoryContents.Where(x => x.Name.EndsWith(".json")))
builder.AddJsonFile(file.PhysicalPath, true, false);
})
.ConfigureServices((hostContext, services) =>
{
services.TryAddKafkaFactories(configuration);
services.AddHostedService<Constellation>();
services.AddHostedService<Qualification>();
});
- Microsoft Azure Eventhub Consumer example: (
ConsumerService.Exampleis example projects)
Constructor Injection in IHostedService
public Constellation(IConsumerFactory factory, ILogger<Constellation> logger)
{
_factory = factory;
_logger = logger;
}
Construct the ICosumerHandle<TKey, TValue> from the factory var handle = _factory.Create<TKey, TValue>(nameof(Constellation));
private IConsumerHandle<TKey, TValue> GetHandle<TKey, TValue>()
{
// Create handle on name registered in Configuration, case sensitive
var handle = _factory.Create<TKey, TValue>(nameof(Constellation));
// Optional Handler Action Setup
handle.Builder
.SetErrorHandler((_, error) => { Log(LogLevel.Error, error.Reason); })
.SetLogHandler((_, message) => { Log(LogLevel.Information, message.Message); });
// Available Handler
// SetStatisticsHandler()
// SetOffsetsCommittedHandler()
// SetPartitionsAssignedHandler()
// SetPartitionsRevokedHandler()
// SetOAuthBearerTokenRefreshHandler()
// Available Key and Value Deserializer Setup
// SetKeyDeserializer()
// SetValueDeserializer()
return handle;
}
Note: A
CustomConsumerBuilderobject is instantiated. Both Handler Actions and Key/Value Deserializer Setup is available.
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
new Thread(() => StartConsumerLoop<string, string>(stoppingToken)).Start();
return Task.CompletedTask;
}
private void StartConsumerLoop<TKey, TValue>(CancellationToken cancellationToken)
{
var handle = GetHandle<TKey, TValue>();
while (!cancellationToken.IsCancellationRequested)
{
try
{
var cr = handle.Consume(cancellationToken);
Log(LogLevel.Information, "{0}: {1}", cr.Message.Key, cr.Message.Value);
}
catch (OperationCanceledException)
{
break;
}
catch (ConsumeException e)
{
Log(LogLevel.Error, "Consume error: {0}", e.Error.Reason);
if (e.Error.IsFatal)
break;
}
catch (Exception e)
{
Log(LogLevel.Error, "Unexpected error: {0}", e.Message);
break;
}
}
}
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net9.0 is compatible. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net9.0
- Confluent.Kafka (>= 2.10.0)
- FluentValidation (>= 12.0.0)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 9.0.5)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.