事务性发件箱模式

This commit is contained in:
administrator committed 2025-03-28 22:27:12 +08:00
1 parent 8da853a351
commit f007572e7f
30 files changed
+978 -135

No files matched your search

@@ -0,0 +1,42 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.EventBus.Abstractions;
using System.Diagnostics.CodeAnalysis;
namespace HelloShop.EventBus.Logging
{
public enum DistributedEventStatus { NotPublished, InProgress, Published, PublishedFailed }
public class DistributedEventLog
{
public Guid EventId { get; set; }
public required string EventTypeName { get; set; }
public required DistributedEvent DistributedEvent { get; set; }
public DistributedEventStatus Status { get; set; } = DistributedEventStatus.NotPublished;
public int TimesSent { get; set; }
public DateTimeOffset CreationTime { get; init; } = TimeProvider.System.GetUtcNow();
public required Guid TransactionId { get; set; }
/// <summary>
/// EF Core cannot set navigation properties using a constructor.
/// The constructor can be public, private, or have any other accessibility.
/// </summary>
private DistributedEventLog() { }
[SetsRequiredMembers]
public DistributedEventLog(DistributedEvent @event, Guid transactionId)
{
EventId = @event.Id;
EventTypeName = @event.GetType().Name;
DistributedEvent = @event;
TransactionId = transactionId;
}
}
}
@@ -0,0 +1,20 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using System.Diagnostics.CodeAnalysis;
namespace HelloShop.EventBus.Logging
{
public static class DistributedEventLogExtensions
{
public static void UseDistributedEventLogs(this ModelBuilder builder) => builder.ApplyConfiguration(new EventLogEntityTypeConfiguration());
public static IServiceCollection AddDistributedEventLogs<TContext>([NotNull] this IServiceCollection services) where TContext : DbContext
{
services.AddTransient<IDistributedEventLogService, DistributedEventLogService<TContext>>().AddHostedService<DistributedEventWorker>();
return services;
}
}
}
@@ -0,0 +1,51 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.EventBus.Abstractions;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Storage;
namespace HelloShop.EventBus.Logging
{
public class DistributedEventLogService<TContext>(TContext dbContext) : IDistributedEventLogService where TContext : DbContext
{
public async Task<IEnumerable<DistributedEventLog>> RetrieveEventLogsPendingToPublishAsync(Guid transactionId, CancellationToken cancellationToken = default)
{
return await dbContext.Set<DistributedEventLog>().Where(e => e.TransactionId == transactionId && e.Status == DistributedEventStatus.NotPublished).ToListAsync(cancellationToken: cancellationToken);
}
public async Task<IEnumerable<DistributedEventLog>> RetrieveEventLogsFailedToPublishAsync(CancellationToken cancellationToken = default)
{
return await dbContext.Set<DistributedEventLog>().Where(e => e.Status == DistributedEventStatus.PublishedFailed).ToListAsync(cancellationToken);
}
public async Task SaveEventAsync(DistributedEvent @event, IDbContextTransaction transaction, CancellationToken cancellationToken = default)
{
var eventLog = new DistributedEventLog(@event, transaction.TransactionId);
dbContext.Database.UseTransaction(transaction.GetDbTransaction());
await dbContext.Set<DistributedEventLog>().AddAsync(eventLog, cancellationToken);
await dbContext.SaveChangesAsync(cancellationToken);
}
public async Task UpdateEventStatusAsync(Guid eventId, DistributedEventStatus status, CancellationToken cancellationToken = default)
{
var eventLogEntry = dbContext.Set<DistributedEventLog>().Single(ie => ie.EventId == eventId);
eventLogEntry.Status = status;
if (status == DistributedEventStatus.InProgress)
{
eventLogEntry.TimesSent++;
}
await dbContext.SaveChangesAsync(cancellationToken);
}
public Task MarkEventAsPublishedAsync(Guid eventId, CancellationToken cancellationToken = default) => UpdateEventStatusAsync(eventId, DistributedEventStatus.Published, cancellationToken);
public Task MarkEventAsInProgressAsync(Guid eventId, CancellationToken cancellationToken = default) => UpdateEventStatusAsync(eventId, DistributedEventStatus.InProgress, cancellationToken);
public Task MarkEventAsFailedAsync(Guid eventId, CancellationToken cancellationToken = default) => UpdateEventStatusAsync(eventId, DistributedEventStatus.PublishedFailed, cancellationToken);
}
}
@@ -0,0 +1,54 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.EventBus.Abstractions;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace HelloShop.EventBus.Logging
{
public class DistributedEventWorker(IServiceScopeFactory serviceScopeFactory) : BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
using var scope = serviceScopeFactory.CreateScope();
var eventBus = scope.ServiceProvider.GetRequiredService<IEventBus>();
var eventLogService = scope.ServiceProvider.GetRequiredService<IDistributedEventLogService>();
var logger = scope.ServiceProvider.GetRequiredService<ILogger<DistributedEventWorker>>();
try
{
var failedEventLogs = await eventLogService.RetrieveEventLogsFailedToPublishAsync(stoppingToken);
foreach (var eventLog in failedEventLogs)
{
DistributedEvent @event = eventLog.DistributedEvent;
try
{
await eventLogService.MarkEventAsInProgressAsync(@event.Id, stoppingToken);
await eventBus.PublishAsync(@event, stoppingToken);
await eventLogService.MarkEventAsPublishedAsync(@event.Id, stoppingToken);
}
catch (Exception ex)
{
logger.LogError(ex, "Publish through event bus failed for {EventId}.", @event.Id);
await eventLogService.MarkEventAsFailedAsync(@event.Id, stoppingToken);
}
}
}
catch (Exception ex)
{
logger.LogError(ex, "An error occurred while retrieving failed event logs.");
}
await Task.Delay(TimeSpan.FromSeconds(30), stoppingToken);
}
}
}
}
@@ -0,0 +1,24 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.EventBus.Abstractions;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Metadata.Builders;
using System.Text.Json;
namespace HelloShop.EventBus.Logging
{
public class EventLogEntityTypeConfiguration : IEntityTypeConfiguration<DistributedEventLog>
{
private static readonly JsonSerializerOptions s_jsonSerializerOptions = new(JsonSerializerDefaults.Web);
public void Configure(EntityTypeBuilder<DistributedEventLog> builder)
{
builder.HasKey(x => x.EventId);
builder.Property(x => x.EventTypeName).HasMaxLength(64);
builder.Property(x => x.Status).HasConversion<string>();
builder.Property(x => x.DistributedEvent).HasConversion(v => JsonSerializer.Serialize(v, v.GetType(), s_jsonSerializerOptions), v => JsonSerializer.Deserialize<DistributedEvent>(v, s_jsonSerializerOptions)!);
}
}
}
@@ -0,0 +1,17 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net9.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.EntityFrameworkCore.Relational" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\HelloShop.EventBus.Abstractions\HelloShop.EventBus.Abstractions.csproj" />
</ItemGroup>
</Project>
@@ -0,0 +1,23 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using HelloShop.EventBus.Abstractions;
using Microsoft.EntityFrameworkCore.Storage;
namespace HelloShop.EventBus.Logging
{
public interface IDistributedEventLogService
{
Task<IEnumerable<DistributedEventLog>> RetrieveEventLogsPendingToPublishAsync(Guid transactionId, CancellationToken cancellationToken = default);
Task<IEnumerable<DistributedEventLog>> RetrieveEventLogsFailedToPublishAsync(CancellationToken cancellationToken = default);
Task SaveEventAsync(DistributedEvent @event, IDbContextTransaction transaction, CancellationToken cancellationToken = default);
Task MarkEventAsPublishedAsync(Guid eventId, CancellationToken cancellationToken = default);
Task MarkEventAsInProgressAsync(Guid eventId, CancellationToken cancellationToken = default);
Task MarkEventAsFailedAsync(Guid eventId, CancellationToken cancellationToken = default);
}
}
@@ -0,0 +1,23 @@
// Copyright (c) HelloShop Corporation. All rights reserved.
// See the license file in the project root for more information.
using Microsoft.EntityFrameworkCore;
namespace HelloShop.EventBus.Logging
{
public class ResilientTransaction(DbContext dbContext)
{
public static ResilientTransaction New(DbContext context) => new(context);
public async Task ExecuteAsync(Func<Task> action)
{
var strategy = dbContext.Database.CreateExecutionStrategy();
await strategy.ExecuteAsync(async () =>
{
await using var transaction = await dbContext.Database.BeginTransactionAsync();
await action();
await transaction.CommitAsync();
});
}
}
}
@@ -7,7 +7,7 @@ namespace HelloShop.EventBus.RabbitMQ
{
public string ExchangeName { get; set; } = "event-bus-exchange";
public required string QueueName { get; set; }
public required string QueueName { get; set; } = "event-bus-queue";
public int RetryCount { get; set; } = 10;
}