Просмотр исходного кода

😎升级RabbitMQ.Client版本

zuohuaijun 1 год назад
Родитель
Сommit
11bf7772bf

+ 0 - 6
Admin.NET/Admin.NET.Application/Admin.NET.Application.csproj

@@ -24,12 +24,6 @@
     </Content>
   </ItemGroup>
 
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
-
   <ItemGroup>
     <ProjectReference Include="..\Admin.NET.Core\Admin.NET.Core.csproj" />
     <ProjectReference Include="..\Plugins\Admin.NET.Plugin.ApprovalFlow\Admin.NET.Plugin.ApprovalFlow.csproj" />

+ 1 - 1
Admin.NET/Admin.NET.Core/Admin.NET.Core.csproj

@@ -34,7 +34,7 @@
     <PackageReference Include="NewLife.Redis" Version="6.0.2024.1202" />
     <PackageReference Include="Novell.Directory.Ldap.NETStandard" Version="3.6.0" />
     <PackageReference Include="QRCoder" Version="1.6.0" />
-    <PackageReference Include="RabbitMQ.Client" Version="6.8.1" />
+    <PackageReference Include="RabbitMQ.Client" Version="7.0.0" />
     <PackageReference Include="SixLabors.ImageSharp.Web" Version="3.1.3" />
     <PackageReference Include="SKIT.FlurlHttpClient.Wechat.Api" Version="3.6.0" />
     <PackageReference Include="SKIT.FlurlHttpClient.Wechat.TenpayV3" Version="3.9.0" />

+ 36 - 22
Admin.NET/Admin.NET.Core/EventBus/RabbitMQEventSourceStore.cs

@@ -13,27 +13,27 @@ namespace Admin.NET.Core;
 /// <summary>
 /// RabbitMQ自定义事件源存储器
 /// </summary>
-public class RabbitMQEventSourceStore : IEventSourceStorer
+public class RabbitMQEventSourceStore : IEventSourceStorer, IDisposable
 {
     /// <summary>
     /// 内存通道事件源存储器
     /// </summary>
-    private readonly Channel<IEventSource> _channel;
+    private Channel<IEventSource> _channelEventSource;
 
     /// <summary>
-    /// 通道对象
+    /// 路由键
     /// </summary>
-    private readonly IModel _model;
+    private string _routeKey;
 
     /// <summary>
     /// 连接对象
     /// </summary>
-    private readonly IConnection _connection;
+    private IConnection _connection;
 
     /// <summary>
-    /// 路由键
+    /// 通道对象
     /// </summary>
-    private readonly string _routeKey;
+    private IChannel _channel;
 
     /// <summary>
     /// 构造函数
@@ -43,30 +43,41 @@ public class RabbitMQEventSourceStore : IEventSourceStorer
     /// <param name="capacity">存储器最多能够处理多少消息,超过该容量进入等待写入</param>
     public RabbitMQEventSourceStore(ConnectionFactory factory, string routeKey, int capacity)
     {
-        // 配置通道,设置超出默认容量后进入等待
+        InitEventSourceStore(factory, routeKey, capacity).GetAwaiter().GetResult();
+    }
+
+    /// <summary>
+    /// 初始化事件源存储器
+    /// </summary>
+    /// <param name="factory">连接工厂</param>
+    /// <param name="routeKey">路由键</param>
+    /// <param name="capacity">存储器最多能够处理多少消息,超过该容量进入等待写入</param>
+    private async Task InitEventSourceStore(ConnectionFactory factory, string routeKey, int capacity)
+    {
+        // 配置通道(超出默认容量后进入等待)
         var boundedChannelOptions = new BoundedChannelOptions(capacity)
         {
             FullMode = BoundedChannelFullMode.Wait
         };
-
         // 创建有限容量通道
-        _channel = Channel.CreateBounded<IEventSource>(boundedChannelOptions);
+        _channelEventSource = Channel.CreateBounded<IEventSource>(boundedChannelOptions);
 
         // 创建连接
-        _connection = factory.CreateConnection();
+        _connection = await factory.CreateConnectionAsync();
+        // 路由键名
         _routeKey = routeKey;
 
         // 创建通道
-        _model = _connection.CreateModel();
+        _channel = await _connection.CreateChannelAsync();
 
         // 声明路由队列
-        _model.QueueDeclare(routeKey, false, false, false, null);
+        await _channel.QueueDeclareAsync(routeKey, false, false, false, null);
 
         // 创建消息订阅者
-        var consumer = new EventingBasicConsumer(_model);
+        var consumer = new AsyncEventingBasicConsumer(_channel);
 
         // 订阅消息并写入内存 Channel
-        consumer.Received += (ch, ea) =>
+        consumer.ReceivedAsync += async (ch, ea) =>
         {
             // 读取原始消息
             var stringEventSource = Encoding.UTF8.GetString(ea.Body.ToArray());
@@ -75,14 +86,14 @@ public class RabbitMQEventSourceStore : IEventSourceStorer
             var eventSource = JSON.Deserialize<ChannelEventSource>(stringEventSource);
 
             // 写入内存管道存储器
-            _channel.Writer.WriteAsync(eventSource);
+            await _channelEventSource.Writer.WriteAsync(eventSource);
 
             // 确认该消息已被消费
-            _model.BasicAck(ea.DeliveryTag, false);
+            await _channel.BasicAckAsync(ea.DeliveryTag, false);
         };
 
         // 启动消费者且设置为手动应答消息
-        _model.BasicConsume(routeKey, false, consumer);
+        await _channel.BasicConsumeAsync(routeKey, false, consumer);
     }
 
     /// <summary>
@@ -101,12 +112,15 @@ public class RabbitMQEventSourceStore : IEventSourceStorer
         {
             // 序列化及发布
             var data = Encoding.UTF8.GetBytes(JSON.Serialize(source));
-            _model.BasicPublish("", _routeKey, null, data);
+            var props = new BasicProperties();
+            props.ContentType = "text/plain";
+            props.DeliveryMode = DeliveryModes.Persistent;
+            await _channel.BasicPublishAsync("", _routeKey, false, props, data);
         }
         else
         {
             // 处理动态订阅
-            await _channel.Writer.WriteAsync(eventSource, cancellationToken);
+            await _channelEventSource.Writer.WriteAsync(eventSource, cancellationToken);
         }
     }
 
@@ -117,7 +131,7 @@ public class RabbitMQEventSourceStore : IEventSourceStorer
     /// <returns>事件源对象</returns>
     public async ValueTask<IEventSource> ReadAsync(CancellationToken cancellationToken)
     {
-        var eventSource = await _channel.Reader.ReadAsync(cancellationToken);
+        var eventSource = await _channelEventSource.Reader.ReadAsync(cancellationToken);
         return eventSource;
     }
 
@@ -126,7 +140,7 @@ public class RabbitMQEventSourceStore : IEventSourceStorer
     /// </summary>
     public void Dispose()
     {
-        _model.Dispose();
+        _channel.Dispose();
         _connection.Dispose();
     }
 }

+ 0 - 3
Admin.NET/Admin.NET.Web.Core/Admin.NET.Web.Core.csproj

@@ -10,9 +10,6 @@
   </PropertyGroup>
 
   <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
     <PackageReference Include="IGeekFan.AspNetCore.Knife4jUI" Version="0.0.16" />
     <PackageReference Include="System.Security.Cryptography.Pkcs" Version="9.0.0" />
   </ItemGroup>

+ 0 - 6
Admin.NET/Admin.NET.Web.Entry/Admin.NET.Web.Entry.csproj

@@ -47,12 +47,6 @@
       <CopyToOutputDirectory>Never</CopyToOutputDirectory>
     </EmbeddedResource>
   </ItemGroup>
-
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
 	
   <ItemGroup>
     <Content Update="wwwroot\upload\logo.png">

+ 0 - 6
Admin.NET/Plugins/Admin.NET.Plugin.ApprovalFlow/Admin.NET.Plugin.ApprovalFlow.csproj

@@ -19,12 +19,6 @@
     </Content>
   </ItemGroup>
 
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
-
   <ItemGroup>
     <ProjectReference Include="..\..\Admin.NET.Core\Admin.NET.Core.csproj" />
   </ItemGroup>

+ 0 - 6
Admin.NET/Plugins/Admin.NET.Plugin.DingTalk/Admin.NET.Plugin.DingTalk.csproj

@@ -10,12 +10,6 @@
     <Description>Admin.NET 通用权限开发平台</Description>
   </PropertyGroup>
 
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
-
   <ItemGroup>
     <None Update="Configuration\**">
       <CopyToOutputDirectory>PreserveNewest</CopyToOutputDirectory>

+ 0 - 6
Admin.NET/Plugins/Admin.NET.Plugin.GoView/Admin.NET.Plugin.GoView.csproj

@@ -10,12 +10,6 @@
     <Description>Admin.NET 通用权限开发平台</Description>
   </PropertyGroup>
 
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
-
   <ItemGroup>
     <None Update="Configuration\**">
       <CopyToOutputDirectory>PreserveNewest</CopyToOutputDirectory>

+ 0 - 6
Admin.NET/Plugins/Admin.NET.Plugin.K3Cloud/Admin.NET.Plugin.K3Cloud.csproj

@@ -10,12 +10,6 @@
     <Description>Admin.NET 通用权限开发平台</Description>
   </PropertyGroup>
 
-  <ItemGroup>
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
-  </ItemGroup>
-
   <ItemGroup>
     <None Update="Configuration\**">
       <CopyToOutputDirectory>PreserveNewest</CopyToOutputDirectory>

+ 0 - 3
Admin.NET/Plugins/Admin.NET.Plugin.ReZero/Admin.NET.Plugin.ReZero.csproj

@@ -25,9 +25,6 @@
 
   <ItemGroup>
     <PackageReference Include="DocumentFormat.OpenXml" Version="3.2.0" />
-    <PackageReference Include="Furion.Extras.Authentication.JwtBearer" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Extras.ObjectMapper.Mapster" Version="4.9.5.26" />
-    <PackageReference Include="Furion.Pure" Version="4.9.5.26" />
     <PackageReference Include="Microsoft.CodeAnalysis.CSharp.Scripting" Version="4.12.0" />
     <PackageReference Include="Rezero.Api" Version="1.7.12" />
   </ItemGroup>