Creos.KafkaHelper 2.0.6

dotnet add package Creos.KafkaHelper --version 2.0.6
                    
NuGet\Install-Package Creos.KafkaHelper -Version 2.0.6
                    
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="Creos.KafkaHelper" Version="2.0.6" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Creos.KafkaHelper" Version="2.0.6" />
                    
Directory.Packages.props
<PackageReference Include="Creos.KafkaHelper" />
                    
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 Creos.KafkaHelper --version 2.0.6
                    
#r "nuget: Creos.KafkaHelper, 2.0.6"
                    
#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 Creos.KafkaHelper@2.0.6
                    
#: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=Creos.KafkaHelper&version=2.0.6
                    
Install as a Cake Addin
#tool nuget:?package=Creos.KafkaHelper&version=2.0.6
                    
Install as a Cake Tool

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:

  1. Pass in a reference to IServiceProvider
  2. Use the IServiceProvider implemenation to get your applicable consumer instance via the Consumer:Name property defined in your Configuration
  3. 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 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. 
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
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