using System; using System.Threading; using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.Logging; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xExceptions.Constants; namespace xEventService.Providers { /// /// an Abstract Kafka Consumer Service ... /// public abstract class XKafkaConsumerServiceBase : XBaseEventConsumerBackgroundService { /// /// Kafka Consumer Class ... /// protected readonly IConsumer consumer; protected XKafkaConsumerServiceBase( string[] allowedActions, string[] allowedSenders, IServiceProvider serviceProvider, XEventServiceConfiguration configuration, ILogger logger, string topic = XEventServiceConstants.XDefaultTopic ) : base( topic: topic, logger: logger, configuration: configuration, allowedActions: allowedActions, allowedSenders: allowedSenders, serviceProvider: serviceProvider ) { // // Validating Configuration ... if (configuration.Broker != XEventBroker.XKafka) { XException.InvalidConfiguration.Throw(); } // // Prepare Kafka Consumer Configuration ... var config = new ConsumerConfig { EnableAutoCommit = false, BootstrapServers = configuration.Url, AutoOffsetReset = AutoOffsetReset.Earliest, GroupId = XEventServiceConstants.XDefaultConsumerGroup, }; consumer = new ConsumerBuilder(config).Build(); } /// /// Default Action Execution for Background Services ... /// /// /// protected override async Task ExecuteAsync( CancellationToken cancellationToken = default ) { // // For Fix Blocking Synchronous ... await Task.Yield(); using var cts = new CancellationTokenSource(); // // Subscribe to Topic ... consumer.Subscribe(topic); // // Consuming ... while (!cancellationToken.IsCancellationRequested) { try { // // Consuming Topic ... var consumeResult = consumer.Consume(cancellationToken); if (!consumeResult.IsNullOrDefault() && !consumeResult.Message.IsNullOrDefault() ) { // await DoWorkAsync( payload: consumeResult.Message, cancellationToken: cancellationToken ); } } catch (ConsumeException) { await Task.Delay(1000, cancellationToken); } catch (Exception) { await Task.Delay(1000, cancellationToken); } } } /// /// Dispose Implementation ... /// public override void Dispose() { // consumer.Close(); consumer.Dispose(); GC.SuppressFinalize(this); base.Dispose(); } } }