Hba.KafkaSdk 1.0.2

dotnet add package Hba.KafkaSdk --version 1.0.2
                    
NuGet\Install-Package Hba.KafkaSdk -Version 1.0.2
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Hba.KafkaSdk" Version="1.0.2" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Hba.KafkaSdk" Version="1.0.2" />
                    
Directory.Packages.props
<PackageReference Include="Hba.KafkaSdk" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Hba.KafkaSdk --version 1.0.2
                    
#r "nuget: Hba.KafkaSdk, 1.0.2"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Hba.KafkaSdk@1.0.2
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Hba.KafkaSdk&version=1.0.2
                    
Install as a Cake Addin
#tool nuget:?package=Hba.KafkaSdk&version=1.0.2
                    
Install as a Cake Tool

Hba.KafkaSdk

English | Français


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 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
1.0.2 184 2/5/2026
1.0.1 121 2/5/2026
1.0.0 127 2/5/2026