FluentKafkaConsumer 1.0.0
dotnet add package FluentKafkaConsumer --version 1.0.0
NuGet\Install-Package FluentKafkaConsumer -Version 1.0.0
<PackageReference Include="FluentKafkaConsumer" Version="1.0.0" />
<PackageVersion Include="FluentKafkaConsumer" Version="1.0.0" />
<PackageReference Include="FluentKafkaConsumer" />
paket add FluentKafkaConsumer --version 1.0.0
#r "nuget: FluentKafkaConsumer, 1.0.0"
#:package FluentKafkaConsumer@1.0.0
#addin nuget:?package=FluentKafkaConsumer&version=1.0.0
#tool nuget:?package=FluentKafkaConsumer&version=1.0.0
FluentKafkaConsumer
A flexible, extensible, and production-ready Kafka consumer library for .NET applications.
Supports strong typing, raw message handling, custom deserializers, OpenTelemetry, and robust configuration via dependency injection.
✨ Features
- ✅ Simple API for consuming Kafka messages with automatic deserialization
- ✅ Support for raw byte message processing
- ✅ Pluggable deserialization formats (JSON, MessagePack, or custom)
- ✅ OpenTelemetry tracing and metrics built-in
- ✅ Fully configurable via DI and
KafkaConsumerSettings - ✅ Built-in error handling and graceful cancellation support
📦 Installation
Install
Install via NuGet:
dotnet add package FluentKafkaConsumer
Or via the NuGet Package Manager:
Install-Package FluentKafkaConsumer
🛠️ Usage
Register the Consumer
services.AddFluentKafkaConsumerServices(new KafkaConsumerSettings
{
Hosts = "localhost:9092",
UserName = "kafka-user",
Password = "kafka-password",
DefaultGroupId = "my-group"
});
Consume Typed Messages
await kafkaConsumerService.ConsumeAsync<MyEvent>(
topic: "my-topic",
logic: async (message, key, timestamp, ct) =>
{
Console.WriteLine($"Received: {message.Id} at {timestamp}");
await Task.CompletedTask;
});
Consume Raw Messages
await kafkaConsumerService.ConsumeAsync(
topic: "binary-topic",
logic: (rawBytes, key, timestamp, ct) =>
{
var raw = Encoding.UTF8.GetString(rawBytes.Span);
Console.WriteLine($"Raw: {raw}");
return Task.CompletedTask;
});
⚙️ KafkaConsumerSettings
public sealed record KafkaConsumerSettings
{
public string Hosts { get; set; }
public string UserName { get; set; }
public string Password { get; set; }
public string DefaultGroupId { get; set; }
public bool EnableAutoCommit { get; set; } = false;
public AutoOffsetReset AutoOffsetReset { get; set; } = AutoOffsetReset.Earliest;
public SecurityProtocol SecurityProtocol { get; set; } = SecurityProtocol.SaslPlaintext;
public SaslMechanism SaslMechanism { get; set; } = SaslMechanism.Plain;
public string ClientId { get; set; } = Dns.GetHostName();
public GroupProtocol? GroupProtocol { get; set; } = null;
public bool EnableMetricsPush { get; set; } = true;
public PartitionAssignmentStrategy PartitionAssignmentStrategy { get; set; } =
PartitionAssignmentStrategy.CooperativeSticky;
}
🔄 Custom Deserialization
To register your own message format, implement IKafkaMessageDeserializer:
public class ProtobufKafkaDeserializer : IKafkaMessageDeserializer
{
public string FormatKey => "protobuf";
public T Deserialize<T>(ReadOnlyMemory<byte> data) where T : class
{
// Your Protobuf logic here
}
}
Then register it:
services.AddSingleton<IKafkaMessageDeserializer, ProtobufKafkaDeserializer>();
Use it when consuming:
await kafkaConsumerService.ConsumeAsync<MyProtoMessage>("proto-topic", handler, formatKey: "protobuf");
📊 Observability
builder.Services.AddOpenTelemetry()
.WithTracing(b => b.AddKafkaConsumerInstrumentation())
.WithMetrics(b => b.AddKafkaConsumerInstrumentation());
📌 Notes
Message deserialization is format-key based, defaulting to JSON.
Deserializers are resolved by key; you can plug in any format.
OpenTelemetry integration is optional but strongly recommended.
📜 License
MIT License
| 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)
- MessagePack (>= 3.1.3)
- OpenTelemetry.Api.ProviderBuilderExtensions (>= 1.12.0)
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.0 | 200 | 5/9/2025 |