ASP.NET Core微服务之开源事件总线CAP的初步使用

释放双眼,带上耳机,听听看~!

Tip: 此篇已加入.NET Core微服务基础系列文章索引

一、CAP简介

*下面的文字来自CAP的Wiki文档:*https://github.com/dotnetcore/CAP/wiki

CAP 是一个在分布式系统中(SOA,MicroService)实现事件总线及最终一致性(分布式事务)的一个开源的 C# 库,她具有轻量级,高性能,易使用等特点。我们可以轻松的在基于 .NET Core 技术的分布式系统中引入CAP,包括但限于 ASP.NET Core 和 ASP.NET Core on .NET Framework。

CAP 的应用场景主要有以下两个:

  • 分布式事务中的最终一致性(异步确保)的方案
  • 具有高可用性的 EventBus

CAP 同时支持使用 RabbitMQ 或 Kafka 进行底层之间的消息发送,我们不需要具备 RabbitMQ 或者 Kafka 的使用经验,仍然可以轻松的将CAP集成到项目中。

CAP 目前支持使用 Sql Server,MySql,PostgreSql 数据库的项目;

CAP 同时支持使用 EntityFrameworkCore 和 Dapper 的项目,可以根据需要选择不同的配置方式;

CAP的作者为园友savorboard(杨晓东),成都地区的.NET社区领导者,棒棒哒!

二、案例结构

此次试验仍然和上一篇基于MassTransit的案例一样(其实是我懒得再改,直接拿来复用),共有四个MicroService应用程序,当用户下订单时会通过CAP作为事件总线发布消息,作为订阅者的库存和配送服务会接收到消息并消费消息。此次试验会采用RabbitMQ作为消息队列,采用MSSQL作为关系型数据库(同时CAP也是支持MSSQL的)。

准备工作:为所有服务通过NuGet安装CAP及其相关包

PM> Install-Package DotNetCore.CAP
下面是RabbitMQ的支持包
PM> Install-Package DotNetCore.CAP.RabbitMQ
下面是MSSQL的支持包
PM> Install-Package DotNetCore.CAP.SqlServer

三、具体实现

3.1 OrderService

(1)启动配置:这里主要需要给CAP指定数据库(它会在这个数据库中创建本地消息表Published和Received)以及使用到的消息队列(这里是RabbitMQ)


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
1    public void ConfigureServices(IServiceCollection services)
2    {
3        services.AddMvc();
4
5        // Repository
6        services.AddScoped<IOrderRepository, OrderRepository>();
7
8        // EF DbContext
9        services.AddDbContext<OrderDbContext>();
10
11        // Dapper-ConnString
12        services.AddSingleton(Configuration["DB:OrderDB"]);
13
14        // CAP
15        services.AddCap(x =>
16        {
17            x.UseEntityFramework<OrderDbContext>(); // EF
18
19            x.UseSqlServer(Configuration["DB:OrderDB"]); // SQL Server
20
21            x.UseRabbitMQ(cfg =>
22            {
23                cfg.HostName = Configuration["MQ:Host"];
24                cfg.VirtualHost = Configuration["MQ:VirtualHost"];
25                cfg.Port = Convert.ToInt32(Configuration["MQ:Port"]);
26                cfg.UserName = Configuration["MQ:UserName"];
27                cfg.Password = Configuration["MQ:Password"];
28            }); // RabbitMQ
29
30            // Below settings is just for demo
31            x.FailedRetryCount = 2;
32            x.FailedRetryInterval = 5;
33        });
34
35        ......
36    }
37
38    // This method gets called by the runtime. Use this method to configure the HTTP request pipeline.
39    public void Configure(IApplicationBuilder app, IHostingEnvironment env, IApplicationLifetime lifetime)
40    {
41        ......
42
43        app.UseMvc();
44
45        // CAP
46        app.UseCap();
47
48        ......
49    }
50

(2)Controller:这里会调用Repository去实现业务逻辑和发送消息


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
1    [Route("api/Order")]
2    public class OrderController : Controller
3    {
4        public IOrderRepository OrderRepository { get; }
5
6        public OrderController(IOrderRepository OrderRepository)
7        {
8            this.OrderRepository = OrderRepository;
9        }
10
11        [HttpPost]
12        public string Post([FromBody]OrderDTO orderDTO)
13        {
14            var result = OrderRepository.CreateOrderByDapper(orderDTO).GetAwaiter().GetResult();
15
16            return result ? "Post Order Success" : "Post Order Failed";
17        }
18    }
19

(3)Repository:这里实现了两种方式:EF和Dapper(基于ADO.NET),其中EF方式中不需要传transaction(当CAP检测到 Publish 是在EF事务区域内的时候,将使用当前的事务上下文进行消息的存储),而基于ADO.NET方式中需要传transaction(由于不能获取到事务上下文,所以需要用户手动的传递事务上下文到CAP中)。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
1    public class OrderRepository : IOrderRepository
2    {
3        public OrderDbContext DbContext { get; }
4        public ICapPublisher CapPublisher { get; }
5        public string ConnStr { get; } // For Dapper use
6
7        public OrderRepository(OrderDbContext DbContext, ICapPublisher CapPublisher, string ConnStr)
8        {
9            this.DbContext = DbContext;
10            this.CapPublisher = CapPublisher;
11            this.ConnStr = ConnStr;
12        }
13
14        public async Task<bool> CreateOrderByEF(IOrder order)
15        {
16            using (var trans = DbContext.Database.BeginTransaction())
17            {
18                var orderEntity = new Order()
19                {
20                    ID = GenerateOrderID(),
21                    OrderUserID = order.OrderUserID,
22                    OrderTime = order.OrderTime,
23                    OrderItems = null,
24                    ProductID = order.ProductID // For demo use
25                };
26
27                DbContext.Orders.Add(orderEntity);
28                await DbContext.SaveChangesAsync();
29
30                // When using EF, no need to pass transaction
31                var orderMessage = new OrderMessage()
32                {
33                    ID = orderEntity.ID,
34                    OrderUserID = orderEntity.OrderUserID,
35                    OrderTime = orderEntity.OrderTime,
36                    OrderItems = null,
37                    ProductID = orderEntity.ProductID // For demo use
38                };
39                
40                await CapPublisher.PublishAsync(EventConstants.EVENT_NAME_CREATE_ORDER, orderMessage);
41
42                trans.Commit();
43            }
44
45            return true;
46        }
47
48        public async Task<bool> CreateOrderByDapper(IOrder order)
49        {
50            using (var conn = new SqlConnection(ConnStr))
51            {
52                conn.Open();
53                using (var trans = conn.BeginTransaction())
54                {
55                    // business code here
56                    string sqlCommand = @"INSERT INTO [dbo].[Orders](OrderID, OrderTime, OrderUserID, ProductID)
57                                                                VALUES(@OrderID, @OrderTime, @OrderUserID, @ProductID)";
58
59                    order.ID = GenerateOrderID();
60                    await conn.ExecuteAsync(sqlCommand, param: new
61                    {
62                        OrderID = order.ID,
63                        OrderTime = DateTime.Now,
64                        OrderUserID = order.OrderUserID,
65                        ProductID = order.ProductID
66                    }, transaction: trans);
67
68                    // For Dapper/ADO.NET, need to pass transaction
69                    var orderMessage = new OrderMessage()
70                    {
71                        ID = order.ID,
72                        OrderUserID = order.OrderUserID,
73                        OrderTime = order.OrderTime,
74                        OrderItems = null,
75                        MessageTime = DateTime.Now,
76                        ProductID = order.ProductID // For demo use
77                    };
78
79                    await CapPublisher.PublishAsync(EventConstants.EVENT_NAME_CREATE_ORDER, orderMessage, trans);
80
81                    trans.Commit();
82                }
83            }
84
85            return true;
86        }
87
88        private string GenerateOrderID()
89        {
90            // TODO: Some business logic to generate Order ID
91            return Guid.NewGuid().ToString();
92        }
93
94        private string GenerateEventID()
95        {
96            // TODO: Some business logic to generate Order ID
97            return Guid.NewGuid().ToString();
98        }
99    }
100

这里摘抄一段CAP wiki中关于事务的一段介绍:

事务在 CAP 具有重要作用,它是保证消息可靠性的一个基石。 在发送一条消息到消息队列的过程中,如果不使用事务,我们是没有办法保证我们的业务代码在执行成功后消息已经成功的发送到了消息队列,或者是消息成功的发送到了消息队列,但是业务代码确执行失败。

这里的失败原因可能是多种多样的,比如连接异常,网络故障等等。

只有业务代码和CAP的Publish代码必须在同一个事务中,才能够保证业务代码和消息代码同时成功或者失败___。_

换句话说,CAP会确保我们这段逻辑中业务代码和消息代码都成功了,才会真正让事务commit。

3.2 StorageService

(1)启动配置:这里主要是指定Subscriber


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
1    public void ConfigureServices(IServiceCollection services)
2    {
3        services.AddMvc();
4
5        // EF DbContext
6        services.AddDbContext<StorageDbContext>();
7
8        // Dapper-ConnString
9        services.AddSingleton(Configuration["DB:StorageDB"]);
10
11        // Subscriber
12        services.AddTransient<IOrderSubscriberService, OrderSubscriberService>();
13
14        // CAP
15        services.AddCap(x =>
16        {
17            x.UseEntityFramework<StorageDbContext>(); // EF
18
19            x.UseSqlServer(Configuration["DB:StorageDB"]); // SQL Server
20
21            x.UseRabbitMQ(cfg =>
22            {
23                cfg.HostName = Configuration["MQ:Host"];
24                cfg.VirtualHost = Configuration["MQ:VirtualHost"];
25                cfg.Port = Convert.ToInt32(Configuration["MQ:Port"]);
26                cfg.UserName = Configuration["MQ:UserName"];
27                cfg.Password = Configuration["MQ:Password"];
28            }); // RabbitMQ
29
30            // Below settings is just for demo
31            x.FailedRetryCount = 2;
32            x.FailedRetryInterval = 5;
33        });
34
35        ......
36    }
37
38    // This method gets called by the runtime. Use this method to configure the HTTP request pipeline.
39    public void Configure(IApplicationBuilder app, IServiceProvider serviceProvider, IHostingEnvironment env, IApplicationLifetime lifetime)
40    {
41        ......
42
43        app.UseMvc();
44
45        // CAP
46        app.UseCap();
47
48        ......
49    }
50

(2)实现Subscriber

首先定义一个接口,建议放到公共类库中


1
2
3
4
5
1    public interface IOrderSubscriberService
2    {
3        Task ConsumeOrderMessage(OrderMessage message);
4    }
5

然后实现这个接口,记得让其实现ICapSubscribe接口,然后我们就可以使用 CapSubscribeAttribute 来订阅 CAP 发布出来的消息。


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
1    public class OrderSubscriberService : IOrderSubscriberService, ICapSubscribe
2    {
3        private readonly string _connStr;
4        
5        public OrderSubscriberService(string connStr)
6        {
7            _connStr = connStr;
8        }
9
10        [CapSubscribe(EventConstants.EVENT_NAME_CREATE_ORDER)]
11        public async Task ConsumeOrderMessage(OrderMessage message)
12        {
13            await Console.Out.WriteLineAsync($"[StorageService] Received message : {JsonHelper.SerializeObject(message)}");
14            await UpdateStorageNumberAsync(message);
15        }
16
17        private async Task<bool> UpdateStorageNumberAsync(OrderMessage order)
18        {
19            //throw new Exception("test"); // just for demo use
20            using (var conn = new SqlConnection(_connStr))
21            {
22                string sqlCommand = @"UPDATE [dbo].[Storages] SET StorageNumber = StorageNumber - 1
23                                                                WHERE StorageID = @ProductID";
24
25                int count = await conn.ExecuteAsync(sqlCommand, param: new
26                {
27                    ProductID = order.ProductID
28                });
29
30                return count > 0;
31            }
32        }
33    }
34

*.CAP约定消息端在方法实现的过程中需要实现幂等性,所谓幂等性就是指用户对于同一操作发起的一次请求或者多次请求的结果是一致的,不会因为多次点击而产生了副作用。这里我没有考虑,实际中需要首先进行验证,避免二次更新

3.3 DeliveryService

(1)启动配置:与StorageService高度类似,只是使用的不是同一个数据库

(2)实现Subscriber


1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
1    public class OrderSubscriberService : IOrderSubscriberService, ICapSubscribe
2    {
3        private readonly string _connStr;
4
5        public OrderSubscriberService(string connStr)
6        {
7            _connStr = connStr;
8        }
9
10        [CapSubscribe(EventConstants.EVENT_NAME_CREATE_ORDER)]
11        public async Task ConsumeOrderMessage(OrderMessage message)
12        {
13            await Console.Out.WriteLineAsync($"[DeliveryService] Received message : {JsonHelper.SerializeObject(message)}");
14            await AddDeliveryRecordAsync(message);
15        }
16
17        private async Task<bool> AddDeliveryRecordAsync(OrderMessage order)
18        {
19            //throw new Exception("test"); // just for demo use
20            using (var conn = new SqlConnection(_connStr))
21            {
22                string sqlCommand = @"INSERT INTO [dbo].[Deliveries](DeliveryID, OrderID, ProductID, OrderUserID, CreatedTime)
23                                                            VALUES (@DeliveryID, @OrderID, @ProductID, @OrderUserID, @CreatedTime)";
24
25                int count = await conn.ExecuteAsync(sqlCommand, param: new
26                {
27                    DeliveryID = Guid.NewGuid().ToString(),
28                    OrderID = order.ID,
29                    OrderUserID = order.OrderUserID,
30                    ProductID = order.ProductID,
31                    CreatedTime = DateTime.Now
32                });
33
34                return count > 0;
35            }
36        }
37    }
38

3.4 快速测试

(1)启动3个微服务,Check 数据库表状态

首先会看到在各个数据库中均创建了本地消息表,这两个表的含义如下:

Cap.Published:这个表主要是用来存储 CAP 发送到MQ(Message Queue)的客户端消息,也就是说你使用 ICapPublisher 接口 Publish 的消息内容。

Cap.Received:这个表主要是用来存储 CAP 接收到 MQ(Message Queue) 的客户端订阅的消息,也就是使用 CapSubscribe[] 订阅的那些消息。

然后看看各个表的数据,目前只有库存表有数据,因为我们要做的只是更新。

(2)通过Postman发一个Post请求

(3)Check控制台输出的日志信息

(4)Check数据库中的业务表和消息表数据:可以看到发送者和接收者都执行成功了,如果其中任何一个参与者发生了异常或者连接不上,CAP会有默认的重试机制(默认是50次最大重试次数,每次重试间隔60s),当失败总次数达到默认失败总次数后,就不会进行重试了,我们可以在 Dashboard 中查看消息失败的原因,然后进行人工重试处理。

另外,由于CAP会在数据库中创建消息表,因此难免会考虑到其性能。CAP提供了一个数据清理的机制,默认情况下会每隔一个小时将消息表的数据进行清理删除,避免数据量过多导致性能的降低。清理规则为 ExpiresAt (字段名)不为空并且小于当前时间的数据。

四、小结

本篇首先简单介绍了一下CAP这个开源项目,然后基于上一篇中的下订单的小案例来进行了基于CAP的改造,并通过一个实例的运行来看到了结果。当然,这个实例并不完美,很多点都没有考虑(比如消息端消费时的幂等性)和失败重试的场景实践等等等等。由于时间和精力的关系,目前只使用到这儿,以后有机会能够应用上会研究下CAP的源码,最后感谢杨晓东为.NET社区带来了一个优秀的开源项目!

示例代码

Click Here => 点我点我

参考资料

CAP – GitHub : https://github.com/dotnetcore/CAP

CAP – Wiki : https://github.com/dotnetcore/CAP/wiki

杨晓东,《BASE:一种ACID的替代方案》

给TA打赏
共{{data.count}}人
人已打赏
安全经验

如何避免Adsense违规封号

2021-10-11 16:36:11

安全经验

安全咨询服务

2022-1-12 14:11:49

个人中心
购物车
优惠劵
今日签到
有新私信 私信列表
搜索