This commit is contained in:
2026-08-18 00:50:33 +03:30
parent 699fc721f5
commit 3e0bb2c25f
5 changed files with 184 additions and 7 deletions
+1
View File
@@ -9,5 +9,6 @@ namespace xEventService.Constants
//
// Defaults ...
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;
namespace xEventService.Interfaces
@@ -15,15 +19,14 @@ namespace xEventService.Interfaces
/// <returns></returns>
Task<bool> SendMessageAsync(
XEventRequest request,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default
);
/// <summary>
/// Sending a Message ...
/// </summary>
/// <param name="message"></param>
/// <param name="topic"></param>
/// <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>
Task<bool> SendMessageAsync(
+2
View File
@@ -1,3 +1,5 @@
using xEventService.Constants;
namespace xEventService.Models
{
/// <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.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<Null, string> producer;
private readonly ILogger<RpkKafkaProducerService> logger;
private readonly ILogger<XKafkaProducerService> logger;
private readonly XEventServiceConfiguration configuration;
public XKafkaProducerService(
ILoggeer<XKafkaProducerService> logger,
ILogger<XKafkaProducerService> 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<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);
}
}
}