Compare commits

...
1 Commits
Author SHA1 Message Date
saherelm 3e0bb2c25f last ... 2026-08-18 00:50:33 +03:30
5 changed files with 184 additions and 7 deletions
+1
View File
@@ -9,5 +9,6 @@ namespace xEventService.Constants
// //
// Defaults ... // Defaults ...
public const string XDefaultTopic = "XSaherelm"; public const string XDefaultTopic = "XSaherelm";
public const string XDefaultConsumerGroup = "XSaherElmGroup";
} }
} }
+6 -3
View File
@@ -1,3 +1,7 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Models; using xEventService.Models;
namespace xEventService.Interfaces namespace xEventService.Interfaces
@@ -15,15 +19,14 @@ namespace xEventService.Interfaces
/// <returns></returns> /// <returns></returns>
Task<bool> SendMessageAsync( Task<bool> SendMessageAsync(
XEventRequest request, XEventRequest request,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default CancellationToken cancellationToken = default
); );
/// <summary> /// <summary>
/// Sending a Message ... /// Sending a Message ...
/// </summary> /// </summary>
/// <param name="message"></param> /// <param name="message">an Instance of <see cref="XEventMessage"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="topic"></param> /// <param name="topic">Topic for Sending Message ...</param>
/// <param name="cancellationToken"></param> /// <param name="cancellationToken"></param>
/// <returns></returns> /// <returns></returns>
Task<bool> SendMessageAsync( Task<bool> SendMessageAsync(
+2
View File
@@ -1,3 +1,5 @@
using xEventService.Constants;
namespace xEventService.Models namespace xEventService.Models
{ {
/// <summary> /// <summary>
+64
View File
@@ -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
{
/// <summary>
/// an Abstract Kafka Consumer Service ...
/// </summary>
public abstract class XKafkaConsumerServiceBase : BackgroundService
{
/// <summary>
/// Consumed Topic Messages ...
/// </summary>
private readonly string topic;
/// <summary>
/// Kafka Consumer Class ...
/// </summary>
public readonly IConsumer<Ignore, string> consumer;
private readonly XEventServiceConfiguration configuration;
private readonly ILogger<XKafkaConsumerServiceBase> logger;
protected XKafkaConsumerServiceBase(
XEventServiceConfiguration configuration,
ILogger<XKafkaConsumerServiceBase> 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<Ignore, string>(config).Build();
}
/// <summary>
/// Dispose Implementation ...
/// </summary>
public override void Dispose()
{
//
consumer.Close();
consumer.Dispose();
GC.SuppressFinalize(this);
base.Dispose();
}
}
}
+110 -3
View File
@@ -3,30 +3,137 @@ using xExceptions.Constants;
using xEventService.Extensions; using xEventService.Extensions;
using xEventService.Interfaces; using xEventService.Interfaces;
using xEventService.Configuration; 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 namespace xEventService.Providers
{ {
public class XKafkaProducerService : IXKafkaProducerService public class XKafkaProducerService : IXKafkaProducerService
{ {
private readonly IProducer<Null, string> producer; private readonly IProducer<Null, string> producer;
private readonly ILogger<RpkKafkaProducerService> logger; private readonly ILogger<XKafkaProducerService> logger;
private readonly XEventServiceConfiguration configuration;
public XKafkaProducerService( public XKafkaProducerService(
ILoggeer<XKafkaProducerService> logger, ILogger<XKafkaProducerService> logger,
XEventServiceConfiguration configuration XEventServiceConfiguration configuration
) )
{ {
// //
this.logger = logger; this.logger = logger;
this.configuration = configuration;
// //
// Validate Configurations ... // Validate Configurations ...
if (!configuration.IsValid() || if (!configuration.IsValid() ||
configuration.Broker != Constants.XEventBroker.XKafka configuration.Broker != XEventBroker.XKafka
) )
{ {
XException.InvalidConfiguration.Throw(); XException.InvalidConfiguration.Throw();
} }
//
// Prepare Kafka Producer Configuration ...
var config = new ProducerConfig
{
Acks = Acks.All,
BootstrapServers = configuration.Url,
MessageSendMaxRetries = configuration.MaxRetries,
};
producer = new ProducerBuilder<Null, string>(config).Build();
}
/// <summary>
/// Sending Message ...
/// </summary>
/// <param name="request">an Instance of <see cref="XEventRequest"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> 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;
}
/// <summary>
/// Sending a Message ...
/// </summary>
/// <param name="message">an Instance of <see cref="XEventMessage"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="topic">Topic for Sending Message ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> 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<Null, string>
{
Value = message
.ToJSON(camelCase: true)
},
cancellationToken: cancellationToken
);
//
result = true;
}
catch
{
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Required Objects ...
/// </summary>
public void Dispose()
{
//
producer.Dispose();
GC.SuppressFinalize(this);
} }
} }
} }