Compare commits
1
Commits
699fc721f5
...
3e0bb2c25f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3e0bb2c25f |
@@ -8,6 +8,7 @@ namespace xEventService.Constants
|
||||
|
||||
//
|
||||
// Defaults ...
|
||||
public const string XDefaultTopic = "XSaherelm";
|
||||
public const string XDefaultTopic = "XSaherelm";
|
||||
public const string XDefaultConsumerGroup = "XSaherElmGroup";
|
||||
}
|
||||
}
|
||||
@@ -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(
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
using xEventService.Constants;
|
||||
|
||||
namespace xEventService.Models
|
||||
{
|
||||
/// <summary>
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user