Creos.KafkaHelper
2.0.6
dotnet add package Creos.KafkaHelper --version 2.0.6
NuGet\Install-Package Creos.KafkaHelper -Version 2.0.6
<PackageReference Include="Creos.KafkaHelper" Version="2.0.6" />
<PackageVersion Include="Creos.KafkaHelper" Version="2.0.6" />
<PackageReference Include="Creos.KafkaHelper" />
paket add Creos.KafkaHelper --version 2.0.6
#r "nuget: Creos.KafkaHelper, 2.0.6"
#:package Creos.KafkaHelper@2.0.6
#addin nuget:?package=Creos.KafkaHelper&version=2.0.6
#tool nuget:?package=Creos.KafkaHelper&version=2.0.6
Creos.KafkaHelper
Supported Frameworks:
- .NET 6
- .NET 8
- .NET 10
Overall
This is a simple library that eases the producing and consuming from Kafka.
It is designed to be configurable via Configuration and allows consuming and producing to multiple topics.
There are three base properties that are configurable within the Configuration:
"KafkaConfiguration": {
"Brokers": [],
"Producers": [],
"Consumers": []
}
Broker List
Brokers accepts a list of 1 or more brokers, depending on how large your cluster is.
Producer List
The Producers list is an optional list of producers within your appliction.
This is designed such that, you can define a single producer, if you wish, and share those producer settings amongst multiple topics. Or if you wish, you can define a unique producer for each topic.
Consumer List
The Consumers list is an optional list of consumers within your application.
Each consumer can be configured to consume from 1 or more topics.
A topic name can even be wildcarded with an asterisk to support multiple topics matching a pattern.
Multiple consumers are processed concurrently.
Configurable Properties
Producer Properties
- ProducerName: (Required) This must be unique amongst all Producers defined within your Configuration. This is an arbitrary string that you will reference within your application to determine which ProducerSettings you want to Produce to.
- Brokers: (Optional) If this is set, this will override the Broker List defined in the parent configuration -- "KafkaConfiguration:Brokers".
- Topic: (Optional) This is the default topic that a Kafka Message is produced to if not explicitly defined within your application. (Not required)
- Active: (Optional) Boolean (true or false) option where this Producer could be not used. Default: true
- Partitioner: (Optional) Possible settings are defined within the Confluent.Kafka.Partitioner enum: Default: ConsistentRandom
- LingerMS: (Optional) Producer Linger defined in milliseconds. Default: 1000
- BatchSizeBytes: (Optional) Producer BatchSize in bytes. Default: 1000000
- MessageMaxBytes: (Optional) Maximum producer request size in bytes, for example 8388608 (8 MiB). Omit to retain the Confluent default. Coordinate with the broker/topic batch limit and allow for overhead.
- CompressionType: (Optional) Confluent producer codec: None, Gzip, Snappy, Lz4, or Zstd. Omit to retain the Confluent default. Example: "CompressionType": "Lz4".
Consumer Properties
- ConsumerName: (Required) This must be unique amongst all Consumers defined within your Configuration. This is an arbitrary string that you will reference within your application to determine which Consumer you want to consumer from.
- Brokers: (Optional) If this is set, this will override the Broker List defined in the parent configuration -- "KafkaConfiguration:Brokers".
- EnableAutoCommit: (Optional) While Kafka guarantees at-least-once delivery. Setting this to false also empowers the application to guarantee at-least-once successful delivery. Default: true
- BatchOffsetsToCommit: (Optional) Default: 500
- FrequencyToCommitMs: (Optional) This setting is ignored if EnableAutoCommit is set to true. Else, this defines the frequency at which the most recently successful consumed offset (per partition) is committed. Default: 1000
- GroupID: (Required) This setting defines the name of the ConsumerGroup.
- Topics: (Required) This setting defines the topic or topics assigned to this consumer. This does support a wildcarded topic name.
- Active: (Optional) Boolean (true or false) option where this Consumer could be not used. Default: true
- ConsumeFromDate: (Optional) Experimental option that allows for an offset to be reset to a particular date. This setting could potentially be removed in a future update.
- StatisticsIntervalMs: (Optional) Overrides default if explicitly set.
- SessionTimeoutMs: (Optional) Overrides default if explicitly set.
- EnablePartitionEof: (Optional) Overrides default if explicitly set.
- ReconnectBackoffMs: (Optional) Overrides default if explicitly set.
- HeartbeatIntervalMs: (Optional) Overrides default if explicitly set.
- ReconnectBackoffMaxMs: (Optional) Overrides default if explicitly set.
- FetchMaxBytes: (Optional) Overrides default if explicitly set.
- AutoCommitIntervalMs: (Optional) Overrides default if explicitly set.
- QueuedMaxMessagesKbytes: (Optional) Overrides default if explicitly set.
Implementation
Consumers
Start a new class that inherits from ConsumerBackgroundService. This is a traditional .NET BackgroundService Within the ConsumerBackgroundService contructor:
- Pass in a reference to IServiceProvider
- Use the IServiceProvider implemenation to get your applicable consumer instance via the Consumer:Name property defined in your Configuration
- Expose your instance of ConsumerMember for later use.
public class ConsumerHostedService : ConsumerBackgroundService
{
private readonly ILogger<ConsumerHostedService> _logger;
private readonly ConsumerMember _consumerMember;
private CancellationToken _token;
public ConsumerHostedService(ILogger<ConsumerHostedService> logger, IServiceProvider serviceProvider)
{
_logger = logger;
_consumerMember = serviceProvider.GetServices<ConsumerMember>().Where(x => x.ConsumerModel.Active && x.ConsumerModel.Name == "TestConsumer").FirstOrDefault();
}
protected override async Task ExecuteAsync(CancellationToken cancellationToken)
{
_token = cancellationToken;
if (_consumerMember != null)
{
await _consumerMember.RegisterConsumerMemberAsync(cancellationToken);
_consumerMember.ConsumeEvent += ProcessConsumedMessageAsync;
}
}
protected override async Task<bool> ProcessConsumedMessageAsync(ConsumeTriggerEventArgs consumeTriggerEvent)
{
var consumeResult = consumeTriggerEvent.ConsumeResult;
_logger.LogDebug("Topic: {Topic}, offset: {Offset}, TopicPartitionOffset: {TopicPartitionOffset}", consumeResult.Topic, consumeResult.Offset, consumeResult.TopicPartitionOffset);
// Your code here:
// return await ProcessMessage(consumeResult);
}
}
You can have multiple instances of a BackgroundService (ConsumerBackgroundService) to process multiple consumers concurrently.
Producers
This is much simpler to use within your application. Inject an instance of IKafkaProducer into your class and call one of the ProduceMessageToKafka methods.
[HttpPost("ProduceMessage")]
public async Task<IActionResult> ProduceMessage()
{
_logger.LogDebug("Entered ConsumerInfo Controller");
await _producer.ProduceMessageToKafkaAsync("Testing", new Message<string, string>
{
Key = "Test_Message",
Value = "Whatever Value Here"
});
return Ok(_consumerAccessor.GetConsumerListJson());
}
Consumer Accessor
This is a class that enables you to access existing configurations and state of your consumers.
There are two exposed methods:
HasConsumerFailed()
This is a simple method that will return a bool if any active consumer has failed.
GetConsumerListJson()
This method will return an object of the current state of all consumers defined in your configuration. Creos.KafkaHelper.Consumer.ConsumerAccessor is an injected transient. The model returned:
public sealed class ConsumerAccessorModel
{
public bool IsActive { get; internal set; }
public DateTime DateTimeLastCommit { get; internal set; }
public ConsumerModel ConsumerModel { get; internal set; }
}
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net6.0 is compatible. net6.0-android was computed. net6.0-ios was computed. net6.0-maccatalyst was computed. net6.0-macos was computed. net6.0-tvos was computed. net6.0-windows was computed. net7.0 was computed. net7.0-android was computed. net7.0-ios was computed. net7.0-maccatalyst was computed. net7.0-macos was computed. net7.0-tvos was computed. net7.0-windows was computed. net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. 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 is compatible. 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. |
-
net10.0
- Confluent.Kafka (>= 2.13.0)
- Microsoft.Extensions.Configuration.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Configuration.Binder (>= 10.0.3)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.3)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.3)
- System.Configuration.ConfigurationManager (>= 10.0.3)
-
net6.0
- Confluent.Kafka (>= 2.13.0)
- Microsoft.Extensions.Configuration.Abstractions (>= 6.0.0)
- Microsoft.Extensions.Configuration.Binder (>= 6.0.0)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 6.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 6.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 6.0.4)
- System.Configuration.ConfigurationManager (>= 6.0.1)
-
net8.0
- Confluent.Kafka (>= 2.13.0)
- Microsoft.Extensions.Configuration.Abstractions (>= 8.0.0)
- Microsoft.Extensions.Configuration.Binder (>= 8.0.0)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 8.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 8.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 8.0.0)
- System.Configuration.ConfigurationManager (>= 8.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.
| Version | Downloads | Last Updated |
|---|---|---|
| 2.0.6 | 132 | 9/14/2026 |
| 2.0.5 | 136 | 4/24/2026 |
| 2.0.4 | 135 | 2/24/2026 |
| 2.0.3 | 108 | 2/24/2026 |
| 2.0.2 | 113 | 2/24/2026 |
| 2.0.1 | 120 | 2/23/2026 |
| 1.0.12 | 116 | 2/18/2026 |
| 1.0.11 | 9,809 | 7/1/2024 |
| 1.0.10 | 289 | 6/26/2024 |
| 1.0.9 | 2,205 | 6/14/2024 |
| 1.0.8 | 4,302 | 5/2/2024 |
| 1.0.7 | 230 | 4/18/2024 |
| 1.0.6 | 175 | 4/18/2024 |
| 1.0.5 | 206 | 4/11/2024 |
| 1.0.4 | 216 | 4/11/2024 |
| 1.0.3 | 178 | 4/10/2024 |
| 1.0.2 | 13,584 | 3/27/2024 |
| 1.0.1 | 189 | 3/19/2024 |
| 1.0.0 | 192 | 3/19/2024 |