使用分布式事件总线异步通信

This commit is contained in:
administrator committed 2024-09-24 22:53:54 +08:00
1 parent 385b09d482
commit c16d6ece15
26 files changed
+1130 -91

No files matched your search

@@ -2,17 +2,19 @@
// See the license file in the project root for more information.
using AutoMapper;
using HelloShop.OrderingService.DistributedEvents.Events;
using HelloShop.OrderingService.Entities.Buyers;
using HelloShop.OrderingService.Entities.Orders;
using HelloShop.OrderingService.Infrastructure;
using HelloShop.OrderingService.LocalEvents;
using HelloShop.OrderingService.Services;
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
using MediatR;
using Microsoft.EntityFrameworkCore;
namespace HelloShop.OrderingService.Commands.Orders
{
public class CreateOrderCommandHandler(IMediator mediator, OrderingServiceDbContext dbContext, IMapper mapper) : IRequestHandler<CreateOrderCommand, bool>
public class CreateOrderCommandHandler(IMediator mediator, OrderingServiceDbContext dbContext, IMapper mapper, IDistributedEventBus distributedEventBus) : IRequestHandler<CreateOrderCommand, bool>
{
public async Task<bool> Handle(CreateOrderCommand request, CancellationToken cancellationToken)
{
@@ -45,6 +47,8 @@ namespace HelloShop.OrderingService.Commands.Orders
await mediator.Publish(new OrderStartedLocalEvent(order), cancellationToken);
await distributedEventBus.PublishAsync(new OrderStartedDistributedEvent(request.UserId), cancellationToken);
return await Task.FromResult(true);
}
}
@@ -0,0 +1,22 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.OrderingService.DistributedEvents.Events;
using HelloShop.OrderingService.Entities.Orders;
using HelloShop.OrderingService.Infrastructure;
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.EventHandling
{
public class OrderStockConfirmedDistributedEventHandler(OrderingServiceDbContext dbContext) : IDistributedEventHandler<OrderStockConfirmedDistributedEvent>
{
public async Task HandleAsync(OrderStockConfirmedDistributedEvent @event)
{
Order order = await dbContext.Set<Order>().FindAsync(@event.OrderId) ?? throw new Exception($"Order with id {@event.OrderId} not found");
order.OrderStatus = OrderStatus.StockConfirmed;
await dbContext.SaveChangesAsync();
}
}
}
@@ -0,0 +1,23 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.OrderingService.DistributedEvents.Events;
using HelloShop.OrderingService.Entities.Orders;
using HelloShop.OrderingService.Infrastructure;
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.EventHandling
{
public class OrderStockRejectedDistributedEventHandler(OrderingServiceDbContext dbContext) : IDistributedEventHandler<OrderStockRejectedDistributedEvent>
{
public async Task HandleAsync(OrderStockRejectedDistributedEvent @event)
{
Order order = await dbContext.Set<Order>().FindAsync(@event.OrderId) ?? throw new Exception($"Order with id {@event.OrderId} not found");
order.OrderStatus = OrderStatus.Cancelled;
order.Description = "Product out of stock.";
await dbContext.SaveChangesAsync();
}
}
}
@@ -0,0 +1,11 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.Events
{
public record OrderAwaitingValidationDistributedEvent(int OrderId, IEnumerable<OrderStockItem> OrderStockItems) : DistributedEvent;
public record OrderStockItem(int ProductId, int Units);
}
@@ -0,0 +1,9 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.Events
{
public record OrderStartedDistributedEvent(int UserId) : DistributedEvent;
}
@@ -0,0 +1,9 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.Events
{
public record OrderStockConfirmedDistributedEvent(int OrderId) : DistributedEvent;
}
@@ -0,0 +1,9 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
namespace HelloShop.OrderingService.DistributedEvents.Events
{
public record OrderStockRejectedDistributedEvent(int OrderId) : DistributedEvent;
}
@@ -5,6 +5,9 @@ using HelloShop.OrderingService.Behaviors;
using HelloShop.OrderingService.Constants;
using HelloShop.OrderingService.Infrastructure;
using HelloShop.OrderingService.Services;
using HelloShop.OrderingService.Workers;
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
using HelloShop.ServiceDefaults.DistributedEvents.DaprBuildingBlocks;
using HelloShop.ServiceDefaults.Extensions;
using Microsoft.EntityFrameworkCore;
using System.Reflection;
@@ -35,11 +38,21 @@ namespace HelloShop.OrderingService.Extensions
builder.Services.AddModelMapper().AddModelValidator();
builder.Services.AddTransient<ISmsSender, MessageService>().AddTransient<IEmailSender, MessageService>();
builder.AddDaprDistributedEventBus().AddSubscriptionFromAssembly();
builder.Services.Configure<HostOptions>(hostOptions =>
{
hostOptions.BackgroundServiceExceptionBehavior = BackgroundServiceExceptionBehavior.Ignore;
});
builder.Services.AddHostedService<GracePeriodWorker>();
}
public static WebApplication MapApplicationEndpoints(this WebApplication app)
{
app.UseDataSeedingProviders();
app.MapDaprDistributedEventBus();
return app;
}
@@ -0,0 +1,55 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.OrderingService.DistributedEvents.Events;
using HelloShop.OrderingService.Entities.Orders;
using HelloShop.OrderingService.Infrastructure;
using HelloShop.ServiceDefaults.DistributedEvents.Abstractions;
using Microsoft.EntityFrameworkCore;
namespace HelloShop.OrderingService.Workers
{
public class GracePeriodWorker(IServiceScopeFactory serviceScopeFactory, ILogger<GracePeriodWorker> logger) : BackgroundService
{
protected async override Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
if (logger.IsEnabled(LogLevel.Information))
{
logger.LogInformation("{Worker} background task is doing background work.", GetType().Name);
}
using var scope = serviceScopeFactory.CreateAsyncScope();
var dbContext = scope.ServiceProvider.GetRequiredService<OrderingServiceDbContext>();
var distributedEventBus = scope.ServiceProvider.GetRequiredService<IDistributedEventBus>();
DateTimeOffset dateTimeOffset = DateTimeOffset.UtcNow.AddMinutes(-1);
try
{
var orders = await dbContext.Set<Order>().Include(o => o.OrderItems).Where(o => o.OrderStatus == OrderStatus.Submitted && o.OrderDate >= dateTimeOffset).ToListAsync(stoppingToken);
foreach (var order in orders)
{
order.OrderStatus = OrderStatus.AwaitingValidation;
var orderStockList = order.OrderItems.Select(orderItem => new OrderStockItem(orderItem.ProductId, orderItem.Units));
await distributedEventBus.PublishAsync(new OrderAwaitingValidationDistributedEvent(order.Id, orderStockList), stoppingToken);
await dbContext.SaveChangesAsync(stoppingToken);
}
}
catch (Exception ex)
{
logger.LogError(ex, "Error occurred executing {Worker} background task.", GetType().Name);
}
await Task.Delay(TimeSpan.FromSeconds(10), stoppingToken);
}
}
}
}