Common Scenarios

The following scenarios cover the most common patterns you will encounter when integrating NexJob into a production application. Each example is self-contained and shows the minimal code needed to solve the problem, along with notes on the key design decisions.

Send a confirmation email after an order

Problem: Your HTTP endpoint creates an order and needs to send a confirmation email, but you don't want the email delivery to block the HTTP response or cause the request to fail if the email service is temporarily unavailable.

Solution: Enqueue the email job after persisting the order and return immediately. The job runs asynchronously on a worker.

// HTTP endpoint
app.MapPost("/orders", async (OrderInput input, IScheduler scheduler, IDb db, HttpContext ctx) =>
{
    var order = await db.Orders.CreateAsync(input, ct: ctx.RequestAborted);

    await scheduler.EnqueueAsync<OrderConfirmationJob, OrderConfirmationInput>(
        new OrderConfirmationInput(order.Id, order.Email),
        queue: "emails",
        cancellationToken: ctx.RequestAborted);

    return Results.Created($"/orders/{order.Id}", order);
});

// Job
public sealed class OrderConfirmationJob : IJob<OrderConfirmationInput>
{
    private readonly IEmailService _email;

    public OrderConfirmationJob(IEmailService email) => _email = email;

    public async Task ExecuteAsync(OrderConfirmationInput input, CancellationToken ct)
    {
        await _email.SendAsync(
            input.Email,
            "Order Confirmed",
            $"Order {input.OrderId} is confirmed.",
            ct);
    }
}

public sealed record OrderConfirmationInput(Guid OrderId, string Email);

Notes: - Route email jobs to a dedicated "emails" queue so a spike in email volume does not slow down other workloads. - Add [Throttle("email-service", maxConcurrent: 10)] to the job class if your email provider enforces concurrency limits.

Process a large file upload asynchronously

Problem: A user uploads a file that takes minutes to process. You must respond to the HTTP request immediately and process the file in the background.

Solution: Save the file to temporary storage, create an import record, enqueue the processor job, and return a 202 Accepted with the import ID the client can poll.

app.MapPost("/imports", async (IFormFile file, IScheduler scheduler, IDb db) =>
{
    var path = await SaveTempAsync(file);
    var import = await db.Imports.CreateAsync(path, status: "pending");

    await scheduler.EnqueueAsync<ImportProcessorJob, ImportInput>(
        new ImportInput(import.Id, path),
        queue: "imports",
        idempotencyKey: $"import-{import.Id}",
        cancellationToken: CancellationToken.None);

    return Results.Accepted($"/imports/{import.Id}", new { ImportId = import.Id });
});

public sealed class ImportProcessorJob : IJob<ImportInput>
{
    private readonly IDb _db;
    private readonly IImportService _import;

    public ImportProcessorJob(IDb db, IImportService import)
    {
        _db = db;
        _import = import;
    }

    public async Task ExecuteAsync(ImportInput input, CancellationToken ct)
    {
        await _db.Imports.UpdateStatusAsync(input.ImportId, "processing", ct);

        try
        {
            await _import.ProcessAsync(input.FilePath, ct);
            await _db.Imports.UpdateStatusAsync(input.ImportId, "completed", ct);
        }
        catch (Exception ex)
        {
            await _db.Imports.UpdateStatusAsync(input.ImportId, $"failed: {ex.Message}", ct);
            throw; // let NexJob handle retries and dead-lettering
        }
    }
}

public sealed record ImportInput(Guid ImportId, string FilePath);

Notes: - The idempotencyKey prevents duplicate jobs if the client retries the upload request. - For very large files, inject IJobContext and call ReportProgressAsync inside the processing loop so the dashboard shows live progress.

Deliver a webhook with exponential backoff

Problem: You need to deliver a webhook payload to a customer-supplied URL. Delivery may fail due to transient network errors or the customer's server being temporarily down.

Solution: Use [Retry] with exponential backoff and let NexJob handle rescheduling automatically on each failure.

[Retry(5, InitialDelay = "00:00:10", Multiplier = 2.0, MaxDelay = "00:05:00")]
public sealed class WebhookDeliveryJob : IJob<WebhookInput>
{
    private readonly HttpClient _http;

    public WebhookDeliveryJob(HttpClient http) => _http = http;

    public async Task ExecuteAsync(WebhookInput input, CancellationToken ct)
    {
        var response = await _http.PostAsJsonAsync(input.Url, input.Payload, ct);
        response.EnsureSuccessStatusCode(); // throws on 4xx/5xx — triggers retry
    }
}

public sealed record WebhookInput(string Url, object Payload);

Retry schedule (with ±10% jitter applied to each delay):

Attempt Delay before next attempt
1 10 s
2 20 s
3 40 s
4 80 s
5 Dead-lettered

Notes: - Add a [Throttle("webhook-sender", maxConcurrent: 20)] attribute to cap outbound concurrency across all webhook jobs. - Register an IDeadLetterHandler<WebhookDeliveryJob> to notify your team or mark the subscription as suspended when all retries are exhausted.

Scheduled daily cleanup with a recurring job

Problem: You need to delete log entries older than 30 days every night at 2 AM without any manual trigger.

Solution: Register a recurring job using a cron expression.

builder.Services.AddNexJob(options =>
{
    options.AddRecurringJob<CleanupOldLogsJob>(
        id: "cleanup-daily",
        cron: "0 2 * * *");  // runs at 02:00 UTC every day
});

public sealed class CleanupOldLogsJob : IJob
{
    private readonly IDbContext _db;

    public CleanupOldLogsJob(IDbContext db) => _db = db;

    public async Task ExecuteAsync(CancellationToken ct)
    {
        var cutoff = DateTimeOffset.UtcNow.AddDays(-30);
        var deleted = await _db.ExecutionLogs
            .Where(l => l.CreatedAt < cutoff)
            .ExecuteDeleteAsync(ct);

        Console.WriteLine($"Deleted {deleted} old log entries");
    }
}

Notes: - In a multi-node deployment, NexJob uses a distributed lock to ensure the recurring job fires exactly once per scheduled interval — not once per node. - Route this job to a low-priority queue so it does not compete with user-facing work during off-peak processing hours.

Prevent duplicate payment processing with idempotency

Problem: A payment webhook from your provider may be delivered more than once. Charging the customer twice would be a serious bug.

Solution: Pass an idempotencyKey when enqueuing the job and set DuplicatePolicy.RejectAlways so NexJob refuses to create a second job for the same order once the first has completed.

// Webhook endpoint — may be called more than once by the payment provider
app.MapPost("/webhooks/payment", async (PaymentEvent evt, IScheduler scheduler) =>
{
    try
    {
        await scheduler.EnqueueAsync<ProcessPaymentJob, PaymentInput>(
            new PaymentInput(evt.OrderId, evt.Amount),
            idempotencyKey: $"payment-{evt.OrderId}",
            duplicatePolicy: DuplicatePolicy.RejectAlways,
            cancellationToken: CancellationToken.None);
    }
    catch (DuplicateJobException)
    {
        // Already processed — acknowledge so the provider stops retrying
    }

    return Results.Ok();
});

public sealed class ProcessPaymentJob : IJob<PaymentInput>
{
    private readonly IPaymentProcessor _processor;

    public ProcessPaymentJob(IPaymentProcessor processor) => _processor = processor;

    public async Task ExecuteAsync(PaymentInput input, CancellationToken ct)
    {
        await _processor.ProcessAsync(input.OrderId, input.Amount, ct);
    }
}

public sealed record PaymentInput(Guid OrderId, decimal Amount);

Notes: - DuplicatePolicy.RejectAlways raises DuplicateJobException even while the first job is still active. Use DuplicatePolicy.RejectIfActive if you want to allow re-enqueueing after the first job completes. - Consider also adding an idempotency check inside ProcessAsync as a defence-in-depth measure.

Resume a long data migration after a restart

Problem: A multi-hour data migration from a legacy API can be interrupted by worker restarts or pod scaling events. Restarting from the beginning after each interruption wastes hours of work.

Solution: Use IJobContext.SaveCheckpointAsync to persist progress after each batch, and GetCheckpoint<TState>() at startup to resume from the last saved position.

public sealed class DataMigrationJob : IJob<MigrationInput>
{
    private readonly IJobContext _context;
    private readonly ILegacyApiClient _client;
    private readonly IDb _db;

    public DataMigrationJob(IJobContext context, ILegacyApiClient client, IDb db)
    {
        _context = context;
        _client = client;
        _db = db;
    }

    public async Task ExecuteAsync(MigrationInput input, CancellationToken ct)
    {
        // Resume from previous checkpoint, or start fresh
        var checkpoint = _context.GetCheckpoint<MigrationCheckpoint>()
                         ?? new MigrationCheckpoint(LastProcessedKey: null, RecordsMigrated: 0);

        var page = await _client.FetchRecordsAsync(
            input.BatchSize, checkpoint.LastProcessedKey, ct);

        while (page.HasItems)
        {
            await _db.BulkInsertAsync(page.Items, ct);

            checkpoint = new MigrationCheckpoint(
                LastProcessedKey: page.LastKey,
                RecordsMigrated: checkpoint.RecordsMigrated + page.Items.Count);

            await _context.SaveCheckpointAsync(
                state: checkpoint,
                percent: null,
                message: $"Migrated {checkpoint.RecordsMigrated} records so far",
                ct: ct);

            page = await _client.FetchRecordsAsync(
                input.BatchSize, checkpoint.LastProcessedKey, ct);
        }
    }
}

public sealed record MigrationCheckpoint(string? LastProcessedKey, int RecordsMigrated);
public sealed record MigrationInput(int BatchSize);

Notes: - On success, NexJob automatically clears the saved checkpoint — no manual cleanup needed. - Combine with [Throttle("legacy-api", maxConcurrent: 2)] if the legacy system has concurrency restrictions.

High-throughput stream ingestion from Kafka or SQS

Problem: You are consuming a high-velocity event stream (Kafka, SQS, or a message bus) and need to process hundreds or thousands of events per second through NexJob.

Solution: Configure batch acknowledgment and a short polling interval to maximise throughput, and attach the appropriate trigger integration.

var builder = WebApplication.CreateBuilder(args);

// Choose your storage backend
builder.Services.AddNexJobSqlServer(
    builder.Configuration.GetConnectionString("NexJobConnection")!);

builder.Services.AddNexJob(options =>
{
    options.Workers = 30;
    options.PollingInterval = TimeSpan.FromMilliseconds(20);

    // Cut database commit round-trips by 90 %+ with async batch acknowledgment
    options.EnableBatchAcknowledgment = true;
})
.AddKafkaTrigger(opt =>
{
    opt.BootstrapServers = "kafka:9092";
    opt.Topic = "high-volume-events";
    opt.GroupId = "event-processing-group";
});

Notes: - Set Workers to match your processing capacity rather than defaulting to the minimum. Monitor p99 job duration and queue depth to tune the value. - Use the [Retention(PurgeOnSuccess = true)] attribute on high-frequency job classes to delete rows immediately on success and avoid table bloat.