Files
xEventService/Providers/XKafkaProducerService.cs
T
2026-08-19 22:44:39 +03:30

140 lines
4.1 KiB
C#

using xCommons.Extensions;
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 : XBaseEventProducerService, IXKafkaProducerService
{
private readonly IProducer<Null, string> producer;
public XKafkaProducerService(
ILogger<XKafkaProducerService> logger,
XEventServiceConfiguration configuration
) : base(configuration, logger)
{
//
// Validate Configurations ...
if (!configuration.IsValid() ||
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(
topic: request.Topic,
message: request.Message,
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 = !response.IsNullOrDefault();
}
catch (Exception ex)
{
//
logger.LogError($"Kafka Produce Error: {ex.Message} ...");
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Required Objects ...
/// </summary>
public void Dispose()
{
//
if (producer != null)
{
producer.Dispose();
}
//
GC.SuppressFinalize(this);
}
}
}