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