如何配置EF与SQL Server适配FastEndpoints的Job Queue?求示例代码
FastEndpoints Job Queue EF+SQL Server 适配实现
1. 定义EF可识别的Job实体
要让EF正确识别继承IJobStorageRecord的实体,需显式标记实体属性并确保接口所有成员完成映射:
using FastEndpoints; using Microsoft.EntityFrameworkCore; using System.ComponentModel.DataAnnotations; using System.ComponentModel.DataAnnotations.Schema; public class JobRecord : IJobStorageRecord { [Key] [DatabaseGenerated(DatabaseGeneratedOption.Identity)] public string Id { get; set; } = Guid.NewGuid().ToString(); [Required] [Column(TypeName = "nvarchar(255)")] public string JobType { get; set; } = string.Empty; [Required] public string Parameters { get; set; } = string.Empty; [Required] public DateTime ScheduledAt { get; set; } public DateTime? CompletedAt { get; set; } public bool IsCompleted { get; set; } }
2. 配置DbContext
将JobRecord添加到DbContext的DbSet中,确保EF能生成对应的SQL Server表:
public class AppDbContext : DbContext { public AppDbContext(DbContextOptions<AppDbContext> options) : base(options) { } public DbSet<JobRecord> JobRecords { get; set; } protected override void OnModelCreating(ModelBuilder modelBuilder) { modelBuilder.Entity<JobRecord>(e => { e.HasKey(j => j.Id); e.Property(j => j.JobType).HasMaxLength(255); e.Property(j => j.Parameters).IsRequired(); e.Property(j => j.ScheduledAt).IsRequired(); }); } }
3. 实现EF版JobStorage
自定义IJobStorage接口实现,用EF操作SQL Server完成任务存储逻辑:
using FastEndpoints; using Microsoft.EntityFrameworkCore; public class EfJobStorage : IJobStorage { private readonly AppDbContext _dbContext; public EfJobStorage(AppDbContext dbContext) { _dbContext = dbContext; } public async Task AddAsync(IJobStorageRecord job, CancellationToken ct = default) { if (job is JobRecord efJob) { await _dbContext.JobRecords.AddAsync(efJob, ct); await _dbContext.SaveChangesAsync(ct); } } public async Task DeleteCompletedAsync(DateTime olderThan, CancellationToken ct = default) { var completedJobs = await _dbContext.JobRecords .Where(j => j.IsCompleted && j.CompletedAt <= olderThan) .ToListAsync(ct); _dbContext.JobRecords.RemoveRange(completedJobs); await _dbContext.SaveChangesAsync(ct); } public async Task<List<IJobStorageRecord>> GetPendingAsync(CancellationToken ct = default) { return await _dbContext.JobRecords .Where(j => !j.IsCompleted && j.ScheduledAt <= DateTime.UtcNow) .OrderBy(j => j.ScheduledAt) .Cast<IJobStorageRecord>() .ToListAsync(ct); } public async Task MarkAsCompletedAsync(string jobID, CancellationToken ct = default) { var job = await _dbContext.JobRecords.FindAsync(new object[] { jobID }, ct); if (job != null) { job.IsCompleted = true; job.CompletedAt = DateTime.UtcNow; await _dbContext.SaveChangesAsync(ct); } } }
4. 注册服务到DI容器
在Program.cs中替换默认JobStorage,注册EF相关服务:
var builder = WebApplication.CreateBuilder(args); // 注册SQL Server DbContext builder.Services.AddDbContext<AppDbContext>(options => options.UseSqlServer(builder.Configuration.GetConnectionString("DefaultConnection"))); // 替换默认JobStorage为EF实现 builder.Services.AddFastEndpoints() .AddJobQueues(options => { options.JobPollInterval = TimeSpan.FromSeconds(10); // 按需调整轮询间隔 }) .ReplaceSingleton<IJobStorage, EfJobStorage>(); var app = builder.Build(); // 开发环境自动迁移数据库(生产环境建议用脚本执行迁移) using (var scope = app.Services.CreateScope()) { var dbContext = scope.ServiceProvider.GetRequiredService<AppDbContext>(); await dbContext.Database.MigrateAsync(); } app.UseFastEndpoints(); app.Run();
5. 关键注意事项
- 确保appsettings.json中
DefaultConnection连接字符串正确指向SQL Server实例 - 生产环境禁止使用
Database.Migrate(),建议通过EF Core迁移脚本完成数据库初始化 - 可根据业务需求调整任务轮询间隔、实体属性的数据库存储类型
内容的提问来源于stack exchange,提问作者Gholamreza Fathpour
相关产品推荐
相关产品推荐

