Hba.KafkaSdk
1.0.2
dotnet add package Hba.KafkaSdk --version 1.0.2
NuGet\Install-Package Hba.KafkaSdk -Version 1.0.2
<PackageReference Include="Hba.KafkaSdk" Version="1.0.2" />
<PackageVersion Include="Hba.KafkaSdk" Version="1.0.2" />
<PackageReference Include="Hba.KafkaSdk" />
paket add Hba.KafkaSdk --version 1.0.2
#r "nuget: Hba.KafkaSdk, 1.0.2"
#:package Hba.KafkaSdk@1.0.2
#addin nuget:?package=Hba.KafkaSdk&version=1.0.2
#tool nuget:?package=Hba.KafkaSdk&version=1.0.2
Hba.KafkaSdk
English
A .NET 9.0 Kafka SDK for microservices, providing producer, consumer, and transactional outbox pattern support for DDD/CQRS architectures.
Features
- Type-safe Producer - Publish events with automatic envelope wrapping and metadata
- Flexible Consumer - Single-handler and multi-handler patterns with automatic routing
- DDD Support - Aggregate ID, type, and version tracking
- Distributed Tracing - Correlation ID and Causation ID propagation
- Transactional Outbox - Database-backed event publishing for consistency
- Security - SASL/SSL authentication support
- Reliability - Idempotent producer, manual offset commit, retry logic
Installation
dotnet add package Hba.KafkaSdk
Quick Start
Configuration
Add to your appsettings.json:
{
"Kafka": {
"BootstrapServers": "localhost:9092",
"ServiceName": "MyService",
"SecurityProtocol": "Plaintext",
"Producer": {
"Acks": "all",
"EnableIdempotence": true
},
"Consumer": {
"AutoOffsetReset": "Earliest",
"EnableAutoCommit": false
}
}
}
Register Services
// Producer
services.AddKafkaProducer(configuration);
// Single consumer
services.AddKafkaConsumer<OrderCreatedEvent, OrderCreatedHandler>("orders-topic");
// Multiple consumers
services.AddKafkaConsumers(builder => builder
.AddHandler<OrderCreatedEvent, OrderCreatedHandler>("orders-topic")
.AddHandler<PaymentReceivedEvent, PaymentHandler>("payments-topic"));
Publish Events
public class OrderService
{
private readonly IEventProducer _producer;
public OrderService(IEventProducer producer) => _producer = producer;
public async Task CreateOrderAsync(Order order)
{
await _producer.PublishAsync(
topic: "orders-topic",
key: order.Id.ToString(),
payload: new OrderCreatedEvent(order.Id, order.Total),
aggregateId: order.Id.ToString(),
aggregateType: "Order",
version: 1);
}
}
Handle Events
public class OrderCreatedHandler : IEventHandler<OrderCreatedEvent>
{
public async Task HandleAsync(EventEnvelope<OrderCreatedEvent> envelope, CancellationToken ct)
{
var @event = envelope.Payload;
// Process the event
Console.WriteLine($"Order {envelope.AggregateId} created with total: {@event.Total}");
}
}
Transactional Outbox
For guaranteed delivery with database transactions:
// Register outbox processor
services.AddOutboxProcessor<MyOutboxRepository>(options =>
{
options.PollingIntervalMs = 1000;
options.BatchSize = 100;
options.MaxRetries = 3;
});
// Save event to outbox within your transaction
var outboxMessage = OutboxMessage.Create(
topic: "orders-topic",
key: order.Id.ToString(),
payload: new OrderCreatedEvent(order.Id, order.Total),
aggregateId: order.Id.ToString(),
aggregateType: "Order");
await _outboxRepository.AddAsync(outboxMessage);
await _unitOfWork.SaveChangesAsync(); // Commit with your transaction
Configuration Options
| Option | Default | Description |
|---|---|---|
BootstrapServers |
localhost:9092 |
Kafka broker addresses |
ServiceName |
- | Consumer group ID |
SecurityProtocol |
Plaintext |
Plaintext, Ssl, SaslPlaintext, SaslSsl |
SaslMechanism |
- | Plain, ScramSha256, ScramSha512 |
Producer.Acks |
all |
all, leader, none |
Producer.EnableIdempotence |
true |
Exactly-once semantics |
Consumer.AutoOffsetReset |
Earliest |
Earliest, Latest |
Consumer.EnableAutoCommit |
false |
Manual offset commit |
Event Envelope
All events are wrapped in EventEnvelope<T> with metadata:
public record EventEnvelope<T>(
Guid EventId,
string EventType,
string AggregateId,
string AggregateType,
int Version,
DateTime Timestamp,
T Payload,
string? CorrelationId,
string? CausationId,
Dictionary<string, string>? Metadata);
Requirements
- .NET 9.0 or later
- Apache Kafka broker
License
MIT License - Hector ADJAKPA
Français
Un SDK Kafka pour .NET 9.0 destiné aux microservices, offrant le support du producteur, du consommateur et du pattern Transactional Outbox pour les architectures DDD/CQRS.
Fonctionnalités
- Producteur typé - Publication d'événements avec enveloppe automatique et métadonnées
- Consommateur flexible - Patterns single-handler et multi-handler avec routage automatique
- Support DDD - Suivi de l'ID d'agrégat, du type et de la version
- Traçage distribué - Propagation des Correlation ID et Causation ID
- Outbox transactionnel - Publication d'événements garantie via base de données
- Sécurité - Support de l'authentification SASL/SSL
- Fiabilité - Producteur idempotent, commit manuel des offsets, logique de retry
Installation
dotnet add package Hba.KafkaSdk
Démarrage rapide
Configuration
Ajoutez dans votre appsettings.json :
{
"Kafka": {
"BootstrapServers": "localhost:9092",
"ServiceName": "MonService",
"SecurityProtocol": "Plaintext",
"Producer": {
"Acks": "all",
"EnableIdempotence": true
},
"Consumer": {
"AutoOffsetReset": "Earliest",
"EnableAutoCommit": false
}
}
}
Enregistrement des services
// Producteur
services.AddKafkaProducer(configuration);
// Consommateur unique
services.AddKafkaConsumer<OrderCreatedEvent, OrderCreatedHandler>("orders-topic");
// Consommateurs multiples
services.AddKafkaConsumers(builder => builder
.AddHandler<OrderCreatedEvent, OrderCreatedHandler>("orders-topic")
.AddHandler<PaymentReceivedEvent, PaymentHandler>("payments-topic"));
Publier des événements
public class OrderService
{
private readonly IEventProducer _producer;
public OrderService(IEventProducer producer) => _producer = producer;
public async Task CreateOrderAsync(Order order)
{
await _producer.PublishAsync(
topic: "orders-topic",
key: order.Id.ToString(),
payload: new OrderCreatedEvent(order.Id, order.Total),
aggregateId: order.Id.ToString(),
aggregateType: "Order",
version: 1);
}
}
Traiter des événements
public class OrderCreatedHandler : IEventHandler<OrderCreatedEvent>
{
public async Task HandleAsync(EventEnvelope<OrderCreatedEvent> envelope, CancellationToken ct)
{
var @event = envelope.Payload;
// Traiter l'événement
Console.WriteLine($"Commande {envelope.AggregateId} créée avec total : {@event.Total}");
}
}
Outbox transactionnel
Pour une livraison garantie avec les transactions de base de données :
// Enregistrer le processeur d'outbox
services.AddOutboxProcessor<MyOutboxRepository>(options =>
{
options.PollingIntervalMs = 1000;
options.BatchSize = 100;
options.MaxRetries = 3;
});
// Sauvegarder l'événement dans l'outbox avec votre transaction
var outboxMessage = OutboxMessage.Create(
topic: "orders-topic",
key: order.Id.ToString(),
payload: new OrderCreatedEvent(order.Id, order.Total),
aggregateId: order.Id.ToString(),
aggregateType: "Order");
await _outboxRepository.AddAsync(outboxMessage);
await _unitOfWork.SaveChangesAsync(); // Commit avec votre transaction
Options de configuration
| Option | Défaut | Description |
|---|---|---|
BootstrapServers |
localhost:9092 |
Adresses des brokers Kafka |
ServiceName |
- | ID du groupe de consommateurs |
SecurityProtocol |
Plaintext |
Plaintext, Ssl, SaslPlaintext, SaslSsl |
SaslMechanism |
- | Plain, ScramSha256, ScramSha512 |
Producer.Acks |
all |
all, leader, none |
Producer.EnableIdempotence |
true |
Sémantique exactly-once |
Consumer.AutoOffsetReset |
Earliest |
Earliest, Latest |
Consumer.EnableAutoCommit |
false |
Commit manuel des offsets |
Enveloppe d'événement
Tous les événements sont encapsulés dans EventEnvelope<T> avec des métadonnées :
public record EventEnvelope<T>(
Guid EventId,
string EventType,
string AggregateId,
string AggregateType,
int Version,
DateTime Timestamp,
T Payload,
string? CorrelationId,
string? CausationId,
Dictionary<string, string>? Metadata);
Prérequis
- .NET 9.0 ou ultérieur
- Broker Apache Kafka
Licence
Licence MIT - Hector ADJAKPA
| 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.6.1)
- Microsoft.Extensions.Configuration.Abstractions (>= 9.0.0)
- Microsoft.Extensions.Configuration.Binder (>= 9.0.0)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 9.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 9.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 9.0.0)
- System.Text.Json (>= 9.0.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.