using System; using Confluent.Kafka; using System.Threading; using xCommons.Extensions; using xEventService.Models; using xExceptions.Constants; using System.Threading.Tasks; using xEventService.Constants; using xEventService.Extensions; using xEventService.Configuration; using Microsoft.Extensions.Logging; 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); logger.LogInformation($"Kafka Cosume Topic: {topic} ..."); // // Consuming ... while (!cancellationToken.IsCancellationRequested) { try { // // Consuming Topic ... var consumeResult = consumer.Consume(cancellationToken); if (!consumeResult.IsNullOrDefault() && !consumeResult.Message.IsNullOrDefault() ) { // var message = consumeResult .Message .Value .FromJSON(); if (message.IsValid( allowedSenders: allowedSenders, allowedActions: allowedActions) ) { // await DoWorkAsync( payload: message, cancellationToken: cancellationToken ); } } } catch (ConsumeException ex) { // logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ..."); await Task.Delay(1000, cancellationToken); } catch (Exception ex) { // logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ..."); await Task.Delay(1000, cancellationToken); } } } /// /// Dispose Object ... /// public override void Dispose() { // if (consumer != null) { consumer.Close(); consumer.Dispose(); } // GC.SuppressFinalize(this); base.Dispose(); } } }