140 lines
4.1 KiB
C#
140 lines
4.1 KiB
C#
using System;
|
|
using Confluent.Kafka;
|
|
using System.Threading;
|
|
using xCommons.Extensions;
|
|
using xEventService.Models;
|
|
using xExceptions.Constants;
|
|
using System.Threading.Tasks;
|
|
using xEventService.Constants;
|
|
using xEventService.Extensions;
|
|
using xEventService.Interfaces;
|
|
using xEventService.Configuration;
|
|
using Microsoft.Extensions.Logging;
|
|
|
|
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);
|
|
}
|
|
}
|
|
} |