From 3e0bb2c25f9a652ddc3086d3f9902f95647364f1 Mon Sep 17 00:00:00 2001 From: Hadi Khazaee Asl Date: Tue, 18 Aug 2026 00:50:33 +0330 Subject: [PATCH] last ... --- Constants/XEventServiceConstants.cs | 3 +- Interfaces/IXKafkaProducerService.cs | 9 +- Models/XEventRequest.cs | 2 + Providers/XKafkaConsumerServiceBase.cs | 64 ++++++++++++++ Providers/XKafkaProducerService.cs | 113 ++++++++++++++++++++++++- 5 files changed, 184 insertions(+), 7 deletions(-) create mode 100644 Providers/XKafkaConsumerServiceBase.cs diff --git a/Constants/XEventServiceConstants.cs b/Constants/XEventServiceConstants.cs index 497ed1e..538258d 100644 --- a/Constants/XEventServiceConstants.cs +++ b/Constants/XEventServiceConstants.cs @@ -8,6 +8,7 @@ namespace xEventService.Constants // // Defaults ... - public const string XDefaultTopic = "XSaherelm"; + public const string XDefaultTopic = "XSaherelm"; + public const string XDefaultConsumerGroup = "XSaherElmGroup"; } } \ No newline at end of file diff --git a/Interfaces/IXKafkaProducerService.cs b/Interfaces/IXKafkaProducerService.cs index da8e0ad..e6f69b2 100644 --- a/Interfaces/IXKafkaProducerService.cs +++ b/Interfaces/IXKafkaProducerService.cs @@ -1,3 +1,7 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using xEventService.Constants; using xEventService.Models; namespace xEventService.Interfaces @@ -15,15 +19,14 @@ namespace xEventService.Interfaces /// Task SendMessageAsync( XEventRequest request, - string topic = XEventServiceConstants.XDefaultTopic, CancellationToken cancellationToken = default ); /// /// Sending a Message ... /// - /// - /// + /// an Instance of to Provides Kafka Messaging requirement ... + /// Topic for Sending Message ... /// /// Task SendMessageAsync( diff --git a/Models/XEventRequest.cs b/Models/XEventRequest.cs index ba1a5e0..3c6a0bc 100644 --- a/Models/XEventRequest.cs +++ b/Models/XEventRequest.cs @@ -1,3 +1,5 @@ +using xEventService.Constants; + namespace xEventService.Models { /// diff --git a/Providers/XKafkaConsumerServiceBase.cs b/Providers/XKafkaConsumerServiceBase.cs new file mode 100644 index 0000000..307dddb --- /dev/null +++ b/Providers/XKafkaConsumerServiceBase.cs @@ -0,0 +1,64 @@ +using System; +using Confluent.Kafka; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using xEventService.Configuration; +using xEventService.Constants; + +namespace xEventService.Providers +{ + /// + /// an Abstract Kafka Consumer Service ... + /// + public abstract class XKafkaConsumerServiceBase : BackgroundService + { + /// + /// Consumed Topic Messages ... + /// + private readonly string topic; + + /// + /// Kafka Consumer Class ... + /// + public readonly IConsumer consumer; + + private readonly XEventServiceConfiguration configuration; + + private readonly ILogger logger; + + protected XKafkaConsumerServiceBase( + XEventServiceConfiguration configuration, + ILogger logger, + string topic = XEventServiceConstants.XDefaultTopic + ) + { + // + this.topic = topic; + this.logger = logger; + this.configuration = configuration; + + // + // Prepare Kafka Consumer Configuration ... + var config = new ConsumerConfig + { + EnableAutoCommit = false, + BootstrapServers = configuration.Url, + AutoOffsetReset = AutoOffsetReset.Earliest, + GroupId = XEventServiceConstants.XDefaultConsumerGroup, + }; + consumer = new ConsumerBuilder(config).Build(); + } + + /// + /// Dispose Implementation ... + /// + public override void Dispose() + { + // + consumer.Close(); + consumer.Dispose(); + GC.SuppressFinalize(this); + base.Dispose(); + } + } +} \ No newline at end of file diff --git a/Providers/XKafkaProducerService.cs b/Providers/XKafkaProducerService.cs index e92339b..1bf4da0 100644 --- a/Providers/XKafkaProducerService.cs +++ b/Providers/XKafkaProducerService.cs @@ -3,30 +3,137 @@ using xExceptions.Constants; using xEventService.Extensions; using xEventService.Interfaces; using xEventService.Configuration; +using xEventService.Models; +using Confluent.Kafka; +using Microsoft.Extensions.Logging; +using System.Threading.Tasks; +using System.Threading; +using xEventService.Constants; +using System; namespace xEventService.Providers { public class XKafkaProducerService : IXKafkaProducerService { private readonly IProducer producer; - private readonly ILogger logger; + private readonly ILogger logger; + private readonly XEventServiceConfiguration configuration; public XKafkaProducerService( - ILoggeer logger, + ILogger logger, XEventServiceConfiguration configuration ) { // this.logger = logger; + this.configuration = configuration; // // Validate Configurations ... if (!configuration.IsValid() || - configuration.Broker != Constants.XEventBroker.XKafka + configuration.Broker != XEventBroker.XKafka ) { XException.InvalidConfiguration.Throw(); } + + // + // Prepare Kafka Producer Configuration ... + var config = new ProducerConfig + { + Acks = Acks.All, + BootstrapServers = configuration.Url, + MessageSendMaxRetries = configuration.MaxRetries, + }; + producer = new ProducerBuilder(config).Build(); + } + + /// + /// Sending Message ... + /// + /// an Instance of to Provides Kafka Messaging requirement ... + /// + /// + public async Task SendMessageAsync( + XEventRequest request, + CancellationToken cancellationToken = default + ) + { + // + // Validate ... + var result = request.IsValid(); + if (!result) + { + XException.InvalidArgs.Throw(); + } + + // + result = await SendMessageAsync( + message: request.Message, + topic: request.Topic, + cancellationToken: cancellationToken + ); + + // + return result; + } + + /// + /// Sending a Message ... + /// + /// an Instance of to Provides Kafka Messaging requirement ... + /// Topic for Sending Message ... + /// + /// + public async Task SendMessageAsync( + XEventMessage message, + string topic = XEventServiceConstants.XDefaultTopic, + CancellationToken cancellationToken = default + ) + { + // + // Validate ... + var result = message.IsValid() && + !topic.IsNullOrEmpty(); + if (!result) + { + XException.InvalidArgs.Throw(); + } + + // + try + { + // + var response = await producer.ProduceAsync( + topic: topic, + message: new Message + { + Value = message + .ToJSON(camelCase: true) + }, + cancellationToken: cancellationToken + ); + + // + result = true; + } + catch + { + result = false; + } + + // + return result; + } + + /// + /// Dispose Required Objects ... + /// + public void Dispose() + { + // + producer.Dispose(); + GC.SuppressFinalize(this); } } } \ No newline at end of file