From 1245a3e390f39c48ed0b51f502e00f8d77edd6b7 Mon Sep 17 00:00:00 2001 From: Halil ibrahim Kalkan Date: Wed, 25 Jul 2018 18:48:26 +0300 Subject: [PATCH] Created RabbitMq connection and channel pooling draft. --- .../RabbitMQ/RabbitMqBackgroundJobManager.cs | 51 +++++++++++++++++++ .../Volo.Abp.RabbitMQ.csproj | 2 +- .../AbpBackgroundJobsRabbitMqModule.cs | 4 ++ .../Volo/Abp/RabbitMQ/AbpRabbitMqOptions.cs | 12 +++++ .../Volo/Abp/RabbitMQ/ChannelPool.cs | 44 ++++++++++++++++ .../Volo/Abp/RabbitMQ/ConnectionPool.cs | 32 ++++++++++++ .../Volo/Abp/RabbitMQ/IChannelAccessor.cs | 10 ++++ .../Volo/Abp/RabbitMQ/IChannelPool.cs | 7 +++ .../Volo/Abp/RabbitMQ/IConnectionPool.cs | 9 ++++ .../Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs | 11 ++++ .../Volo/Abp/RabbitMQ/RabbitMqConnections.cs | 25 +++++++++ .../RabbitMQ/Utf8JsonRabbitMqSerializer.cs | 32 ++++++++++++ 12 files changed, 238 insertions(+), 1 deletion(-) create mode 100644 framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/RabbitMqBackgroundJobManager.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqOptions.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs create mode 100644 framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs diff --git a/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/RabbitMqBackgroundJobManager.cs b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/RabbitMqBackgroundJobManager.cs new file mode 100644 index 0000000000..9892727d54 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/RabbitMqBackgroundJobManager.cs @@ -0,0 +1,51 @@ +using System; +using System.Threading.Tasks; +using RabbitMQ.Client; +using Volo.Abp.DependencyInjection; +using Volo.Abp.RabbitMQ; + +namespace Volo.Abp.BackgroundJobs.RabbitMQ +{ + public class RabbitMqBackgroundJobManager : IBackgroundJobManager, ITransientDependency + { + protected IChannelPool ChannelPool { get; } + protected IRabbitMqSerializer Serializer { get; } + + public RabbitMqBackgroundJobManager(IChannelPool channelPool, IRabbitMqSerializer serializer) + { + Serializer = serializer; + ChannelPool = channelPool; + } + + public Task EnqueueAsync(TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, + TimeSpan? delay = null) + { + var jobName = BackgroundJobNameAttribute.GetName(); + var queueName = "BackgroundJobs." + jobName; //TODO: Make prefix optional + + using (var channelAccessor = ChannelPool.Acquire(queueName)) + { + var properties = channelAccessor.Channel.CreateBasicProperties(); + properties.Persistent = true; + + Publish(channelAccessor.Channel, queueName, args, properties); + } + + return null; //TODO: Can RabbitMQ return a message identifier? + } + + private void Publish( + IModel channel, + string queueName, + TArgs args, + IBasicProperties properties) + { + channel.BasicPublish( + exchange: "", + routingKey: queueName, + basicProperties: properties, + body: Serializer.Serialize(args) + ); + } + } +} diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo.Abp.RabbitMQ.csproj b/framework/src/Volo.Abp.RabbitMQ/Volo.Abp.RabbitMQ.csproj index 580e38a16b..5173519237 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo.Abp.RabbitMQ.csproj +++ b/framework/src/Volo.Abp.RabbitMQ/Volo.Abp.RabbitMQ.csproj @@ -15,7 +15,7 @@ - + diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs index fb283dc609..01e2e52d06 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpBackgroundJobsRabbitMqModule.cs @@ -1,8 +1,12 @@ using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.Json; using Volo.Abp.Modularity; namespace Volo.Abp.RabbitMQ { + [DependsOn( + typeof(AbpJsonModule) + )] public class AbpRabbitMqModule : AbpModule { public override void ConfigureServices(ServiceConfigurationContext context) diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqOptions.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqOptions.cs new file mode 100644 index 0000000000..a971b85d5e --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqOptions.cs @@ -0,0 +1,12 @@ +namespace Volo.Abp.RabbitMQ +{ + public class AbpRabbitMqOptions + { + public RabbitMqConnections Connections { get; } + + public AbpRabbitMqOptions() + { + Connections = new RabbitMqConnections(); + } + } +} diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs new file mode 100644 index 0000000000..3e3d7ab6ad --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs @@ -0,0 +1,44 @@ +using RabbitMQ.Client; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.RabbitMQ +{ + public class ChannelPool : IChannelPool, ISingletonDependency + { + protected IConnectionPool ConnectionPool { get; } + + public ChannelPool(IConnectionPool connectionPool) + { + ConnectionPool = connectionPool; + } + + public virtual IChannelAccessor Acquire(string channelName = null) + { + //TODO: Pool channels! + return new ChannelAccessor( + CreateChannel(channelName) + ); + } + + protected virtual IModel CreateChannel(string channelName) + { + //TODO: How to determine the right connection name? + return ConnectionPool.Get().CreateModel(); + } + + protected class ChannelAccessor : IChannelAccessor + { + public IModel Channel { get; } + + public ChannelAccessor(IModel channel) + { + Channel = channel; + } + + public void Dispose() + { + Channel.Dispose(); + } + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs new file mode 100644 index 0000000000..a173564d52 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs @@ -0,0 +1,32 @@ +using System.Collections.Concurrent; +using System.Collections.Generic; +using Microsoft.Extensions.Options; +using RabbitMQ.Client; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.RabbitMQ +{ + public class ConnectionPool : IConnectionPool, ISingletonDependency + { + protected AbpRabbitMqOptions Options { get; } + + protected ConcurrentDictionary Connections { get; } + + public ConnectionPool(IOptions options) + { + Options = options.Value; + Connections = new ConcurrentDictionary(); + } + + public virtual IConnection Get(string connectionName = null) + { + return Connections.GetOrAdd(connectionName, () => + { + var connectionFactory = Options.Connections.GetOrDefault(connectionName) + ?? Options.Connections.Default; + + return connectionFactory.CreateConnection(); + }); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs new file mode 100644 index 0000000000..32a6584d62 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelAccessor.cs @@ -0,0 +1,10 @@ +using System; +using RabbitMQ.Client; + +namespace Volo.Abp.RabbitMQ +{ + public interface IChannelAccessor : IDisposable + { + IModel Channel { get; } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs new file mode 100644 index 0000000000..36c011d098 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IChannelPool.cs @@ -0,0 +1,7 @@ +namespace Volo.Abp.RabbitMQ +{ + public interface IChannelPool + { + IChannelAccessor Acquire(string channelName = null); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs new file mode 100644 index 0000000000..8e0b06bebd --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IConnectionPool.cs @@ -0,0 +1,9 @@ +using RabbitMQ.Client; + +namespace Volo.Abp.RabbitMQ +{ + public interface IConnectionPool + { + IConnection Get(string connectionName = null); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs new file mode 100644 index 0000000000..771d1a05a6 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs @@ -0,0 +1,11 @@ +using System; + +namespace Volo.Abp.RabbitMQ +{ + public interface IRabbitMqSerializer + { + byte[] Serialize(object obj); + + object Deserialize(byte[] value, Type type); + } +} diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs new file mode 100644 index 0000000000..9e543c4091 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs @@ -0,0 +1,25 @@ +using System; +using System.Collections.Generic; +using JetBrains.Annotations; +using RabbitMQ.Client; + +namespace Volo.Abp.RabbitMQ +{ + [Serializable] + public class RabbitMqConnections : Dictionary + { + public const string DefaultConnectionName = "Default"; + + [NotNull] + public ConnectionFactory Default + { + get => this.GetOrDefault(DefaultConnectionName); + set => this[DefaultConnectionName] = Check.NotNull(value, nameof(value)); + } + + public RabbitMqConnections() + { + Default = new ConnectionFactory(); + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs new file mode 100644 index 0000000000..f097340b37 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs @@ -0,0 +1,32 @@ +using System; +using System.Text; +using Volo.Abp.DependencyInjection; +using Volo.Abp.Json; + +namespace Volo.Abp.RabbitMQ +{ + public class Utf8JsonRabbitMqSerializer : IRabbitMqSerializer, ITransientDependency + { + private readonly IJsonSerializer _jsonSerializer; + + public Utf8JsonRabbitMqSerializer(IJsonSerializer jsonSerializer) + { + _jsonSerializer = jsonSerializer; + } + + public byte[] Serialize(object obj) + { + return Encoding.UTF8.GetBytes(_jsonSerializer.Serialize(obj)); + } + + public object Deserialize(byte[] value, Type type) + { + return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); + } + + public T Deserialize(string value) + { + return _jsonSerializer.Deserialize(value); + } + } +} \ No newline at end of file