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