You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何配置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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 00:38:11