| name | dotnet-outbox-pattern |
| description | Implements the Outbox pattern for reliable domain event processing. Ensures events are persisted in the same transaction as the aggregate changes and processed asynchronously with guaranteed delivery. |
| version | 1.0.0 |
| language | C# |
| framework | .NET 8+ |
| dependencies | Entity Framework Core, Quartz.NET, MediatR |
| pattern | Transactional Outbox, Guaranteed Delivery |
Outbox Pattern Implementation
Overview
The Outbox pattern ensures reliable event processing:
- Atomic persistence - Events saved in same transaction as aggregate
- Guaranteed delivery - Events processed even if app crashes
- Eventual consistency - Async processing with retry
- Idempotency - Handle duplicate processing gracefully
Quick Reference
| Component | Purpose |
|---|
OutboxMessage | Persisted event entity |
OutboxMessageConfiguration | EF Core mapping |
ProcessOutboxMessagesJob | Background processor (Quartz) |
IdempotentDomainEventHandler | Deduplicated handler wrapper |
OutboxConsumer | Alternative direct DB poller |
Outbox Structure
/Infrastructure/
├── Outbox/
│ ├── OutboxMessage.cs
│ ├── OutboxMessageConfiguration.cs
│ ├── ProcessOutboxMessagesJob.cs
│ ├── ProcessOutboxMessagesJobSetup.cs
│ └── IdempotentDomainEventHandler.cs
└── ApplicationDbContext.cs
Template: Outbox Message Entity
namespace {name}.infrastructure.outbox;
public sealed class OutboxMessage
{
public OutboxMessage()
{
}
public OutboxMessage(Guid id, string type, string content, DateTime occurredOnUtc)
{
Id = id;
Type = type;
Content = content;
OccurredOnUtc = occurredOnUtc;
}
public Guid Id { get; set; }
public string Type { get; set; } = string.Empty;
public string Content { get; set; } = string.Empty;
public DateTime OccurredOnUtc { get; set; }
public DateTime? ProcessedOnUtc { get; set; }
public string? Error { get; set; }
public int RetryCount { get; set; }
}
Template: EF Core Configuration
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Metadata.Builders;
namespace {name}.infrastructure.outbox;
internal sealed class OutboxMessageConfiguration
: IEntityTypeConfiguration<OutboxMessage>
{
public void Configure(EntityTypeBuilder<OutboxMessage> builder)
{
builder.ToTable("outbox_message");
builder.HasKey(o => o.Id);
builder.Property(o => o.Id)
.ValueGeneratedNever();
builder.Property(o => o.Type)
.HasMaxLength(500)
.IsRequired();
builder.Property(o => o.Content)
.HasColumnType("jsonb")
.IsRequired();
builder.Property(o => o.OccurredOnUtc)
.IsRequired();
builder.Property(o => o.ProcessedOnUtc);
builder.Property(o => o.Error)
.HasColumnType("text");
builder.Property(o => o.RetryCount)
.HasDefaultValue(0);
builder.HasIndex(o => o.ProcessedOnUtc)
.HasFilter("processed_on_utc IS NULL")
.HasDatabaseName("ix_outbox_message_unprocessed");
builder.HasIndex(o => o.ProcessedOnUtc)
.HasFilter("processed_on_utc IS NOT NULL")
.HasDatabaseName("ix_outbox_message_processed");
}
}
Template: DbContext Integration
using System.Text.Json;
using Microsoft.EntityFrameworkCore;
using {name}.domain.abstractions;
using {name}.infrastructure.outbox;
namespace {name}.infrastructure;
public sealed class ApplicationDbContext : DbContext, IUnitOfWork
{
private static readonly JsonSerializerOptions JsonOptions = new()
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
WriteIndented = false
};
public ApplicationDbContext(DbContextOptions<ApplicationDbContext> options)
: base(options)
{
}
public DbSet<OutboxMessage> OutboxMessages => Set<OutboxMessage>();
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
modelBuilder.ApplyConfigurationsFromAssembly(typeof(ApplicationDbContext).Assembly);
base.OnModelCreating(modelBuilder);
}
public override async Task<int> SaveChangesAsync(CancellationToken cancellationToken = default)
{
ConvertDomainEventsToOutboxMessages();
return await base.SaveChangesAsync(cancellationToken);
}
private void ConvertDomainEventsToOutboxMessages()
{
var entitiesWithEvents = ChangeTracker
.Entries<Entity>()
.Where(e => e.Entity.GetDomainEvents().Any())
.Select(e => e.Entity)
.ToList();
var domainEvents = entitiesWithEvents
.SelectMany(e => e.GetDomainEvents())
.ToList();
foreach (var entity in entitiesWithEvents)
{
entity.ClearDomainEvents();
}
foreach (var domainEvent in domainEvents)
{
var outboxMessage = new OutboxMessage
{
Id = Guid.NewGuid(),
Type = domainEvent.GetType().AssemblyQualifiedName!,
Content = JsonSerializer.Serialize(
domainEvent,
domainEvent.GetType(),
JsonOptions),
OccurredOnUtc = DateTime.UtcNow
};
OutboxMessages.Add(outboxMessage);
}
}
}
Template: Outbox Processor Job (Quartz)
using System.Text.Json;
using MediatR;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using Quartz;
using {name}.domain.abstractions;
namespace {name}.infrastructure.outbox;
[DisallowConcurrentExecution]
public sealed class ProcessOutboxMessagesJob : IJob
{
private const int BatchSize = 20;
private const int MaxRetries = 3;
private static readonly JsonSerializerOptions JsonOptions = new()
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase
};
private readonly ApplicationDbContext _dbContext;
private readonly IPublisher _publisher;
private readonly ILogger<ProcessOutboxMessagesJob> _logger;
public ProcessOutboxMessagesJob(
ApplicationDbContext dbContext,
IPublisher publisher,
ILogger<ProcessOutboxMessagesJob> logger)
{
_dbContext = dbContext;
_publisher = publisher;
_logger = logger;
}
public async Task Execute(IJobExecutionContext context)
{
_logger.LogDebug("Starting outbox message processing...");
var messages = await GetUnprocessedMessages(context.CancellationToken);
if (!messages.Any())
{
_logger.LogDebug("No outbox messages to process");
return;
}
_logger.LogInformation(
"Processing {Count} outbox messages",
messages.Count);
foreach (var message in messages)
{
await ProcessMessage(message, context.CancellationToken);
}
await _dbContext.SaveChangesAsync(context.CancellationToken);
_logger.LogInformation("Completed outbox message processing");
}
private async Task<List<OutboxMessage>> GetUnprocessedMessages(
CancellationToken cancellationToken)
{
return await _dbContext.OutboxMessages
.Where(m => m.ProcessedOnUtc == null)
.Where(m => m.RetryCount < MaxRetries)
.OrderBy(m => m.OccurredOnUtc)
.Take(BatchSize)
.ToListAsync(cancellationToken);
}
private async Task ProcessMessage(
OutboxMessage message,
CancellationToken cancellationToken)
{
try
{
_logger.LogDebug(
"Processing outbox message {MessageId} of type {Type}",
message.Id,
message.Type);
var eventType = Type.GetType(message.Type);
if (eventType is null)
{
_logger.LogError(
"Could not resolve type {Type} for message {MessageId}",
message.Type,
message.Id);
message.Error = $"Could not resolve type: {message.Type}";
message.ProcessedOnUtc = DateTime.UtcNow;
return;
}
var domainEvent = JsonSerializer.Deserialize(
message.Content,
eventType,
JsonOptions) as IDomainEvent;
if (domainEvent is null)
{
_logger.LogError(
"Could not deserialize message {MessageId}",
message.Id);
message.Error = "Could not deserialize message content";
message.ProcessedOnUtc = DateTime.UtcNow;
return;
}
await _publisher.Publish(domainEvent, cancellationToken);
message.ProcessedOnUtc = DateTime.UtcNow;
message.Error = null;
_logger.LogInformation(
"Successfully processed outbox message {MessageId}",
message.Id);
}
catch (Exception ex)
{
_logger.LogError(
ex,
"Error processing outbox message {MessageId}. Retry count: {RetryCount}",
message.Id,
message.RetryCount);
message.RetryCount++;
message.Error = ex.ToString();
if (message.RetryCount >= MaxRetries)
{
message.ProcessedOnUtc = DateTime.UtcNow;
_logger.LogError(
"Outbox message {MessageId} exceeded max retries and has been marked as failed",
message.Id);
}
}
}
}
Template: Job Configuration
using Microsoft.Extensions.Options;
using Quartz;
namespace {name}.infrastructure.outbox;
internal sealed class ProcessOutboxMessagesJobSetup
: IConfigureOptions<QuartzOptions>
{
public void Configure(QuartzOptions options)
{
var jobKey = JobKey.Create(nameof(ProcessOutboxMessagesJob));
options
.AddJob<ProcessOutboxMessagesJob>(jobBuilder =>
jobBuilder.WithIdentity(jobKey))
.AddTrigger(triggerBuilder =>
triggerBuilder
.ForJob(jobKey)
.WithSimpleSchedule(schedule =>
schedule
.WithIntervalInSeconds(10)
.RepeatForever()));
}
}
Template: Idempotent Event Handler Wrapper
using MediatR;
using Microsoft.EntityFrameworkCore;
using {name}.domain.abstractions;
namespace {name}.infrastructure.outbox;
public abstract class IdempotentDomainEventHandler<TEvent>
: INotificationHandler<TEvent>
where TEvent : IDomainEvent
{
private readonly ApplicationDbContext _dbContext;
protected IdempotentDomainEventHandler(ApplicationDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task Handle(TEvent notification, CancellationToken cancellationToken)
{
var handlerName = GetType().Name;
var eventId = notification.Id;
var alreadyProcessed = await _dbContext
.Set<OutboxMessageConsumer>()
.AnyAsync(
c => c.EventId == eventId && c.HandlerName == handlerName,
cancellationToken);
if (alreadyProcessed)
{
return;
}
await HandleAsync(notification, cancellationToken);
_dbContext.Set<OutboxMessageConsumer>().Add(new OutboxMessageConsumer
{
Id = Guid.NewGuid(),
EventId = eventId,
HandlerName = handlerName,
ProcessedOnUtc = DateTime.UtcNow
});
await _dbContext.SaveChangesAsync(cancellationToken);
}
protected abstract Task HandleAsync(TEvent notification, CancellationToken cancellationToken);
}
public sealed class OutboxMessageConsumer
{
public Guid Id { get; set; }
public Guid EventId { get; set; }
public string HandlerName { get; set; } = string.Empty;
public DateTime ProcessedOnUtc { get; set; }
}
Template: Cleanup Job
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Logging;
using Quartz;
namespace {name}.infrastructure.outbox;
[DisallowConcurrentExecution]
public sealed class CleanupOutboxMessagesJob : IJob
{
private const int RetentionDays = 7;
private const int BatchSize = 1000;
private readonly ApplicationDbContext _dbContext;
private readonly ILogger<CleanupOutboxMessagesJob> _logger;
public CleanupOutboxMessagesJob(
ApplicationDbContext dbContext,
ILogger<CleanupOutboxMessagesJob> logger)
{
_dbContext = dbContext;
_logger = logger;
}
public async Task Execute(IJobExecutionContext context)
{
var cutoffDate = DateTime.UtcNow.AddDays(-RetentionDays);
_logger.LogInformation(
"Cleaning up outbox messages processed before {CutoffDate}",
cutoffDate);
var totalDeleted = 0;
int deletedInBatch;
do
{
deletedInBatch = await _dbContext.OutboxMessages
.Where(m => m.ProcessedOnUtc != null)
.Where(m => m.ProcessedOnUtc < cutoffDate)
.Take(BatchSize)
.ExecuteDeleteAsync(context.CancellationToken);
totalDeleted += deletedInBatch;
} while (deletedInBatch == BatchSize);
_logger.LogInformation(
"Cleaned up {Count} old outbox messages",
totalDeleted);
}
}
Template: Registration
private static void AddBackgroundJobs(
IServiceCollection services,
IConfiguration configuration)
{
services.AddQuartz(configure =>
{
configure.UsePersistentStore(options =>
{
options.UsePostgres(configuration.GetConnectionString("Database")!);
options.UseJsonSerializer();
});
});
services.AddQuartzHostedService(options =>
{
options.WaitForJobsToComplete = true;
});
services.ConfigureOptions<ProcessOutboxMessagesJobSetup>();
services.ConfigureOptions<CleanupOutboxMessagesJobSetup>();
}
Database Migration
CREATE TABLE outbox_message (
id UUID PRIMARY KEY,
type VARCHAR(500) NOT NULL,
content JSONB NOT NULL,
occurred_on_utc TIMESTAMP NOT NULL,
processed_on_utc TIMESTAMP NULL,
error TEXT NULL,
retry_count INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX ix_outbox_message_unprocessed
ON outbox_message (occurred_on_utc)
WHERE processed_on_utc IS NULL;
CREATE INDEX ix_outbox_message_processed
ON outbox_message (processed_on_utc)
WHERE processed_on_utc IS NOT NULL;
CREATE TABLE outbox_message_consumer (
id UUID PRIMARY KEY,
event_id UUID NOT NULL,
handler_name VARCHAR(500) NOT NULL,
processed_on_utc TIMESTAMP NOT NULL
);
CREATE UNIQUE INDEX ix_outbox_consumer_event_handler
ON outbox_message_consumer (event_id, handler_name);
Critical Rules
- Same transaction - Events saved with aggregate in one transaction
- Idempotent handlers - Must handle duplicate delivery
- Order not guaranteed - Events may process out of order
- Retry with backoff - Don't hammer failing events
- Cleanup old messages - Prevent table bloat
- Monitor failures - Alert on max retries exceeded
- Type serialization - Use
AssemblyQualifiedName for deserialize
- JSON serialization - Consistent options for serialize/deserialize
- Batch processing - Don't process one at a time
- Disable concurrent execution - Prevent duplicate processing
Anti-Patterns to Avoid
await _publisher.Publish(new UserCreatedEvent(user.Id));
await _unitOfWork.SaveChangesAsync();
user.RaiseDomainEvent(new UserCreatedEvent(user.Id));
await _unitOfWork.SaveChangesAsync();
public async Task Handle(UserCreatedEvent e, CancellationToken ct)
{
await _emailService.SendWelcomeEmail(e.UserId);
}
public async Task Handle(UserCreatedEvent e, CancellationToken ct)
{
if (await _emailLog.ExistsAsync(e.UserId, "welcome"))
return;
await _emailService.SendWelcomeEmail(e.UserId);
await _emailLog.RecordAsync(e.UserId, "welcome");
}
foreach (var message in allMessages)
{
await ProcessMessage(message);
}
var messages = await _dbContext.OutboxMessages
.Where(m => m.ProcessedOnUtc == null)
.Take(20)
.ToListAsync();
Related Skills
dotnet-domain-events-generator - Domain events that go into outbox
dotnet-quartz-background-jobs - Background job scheduling
dotnet-clean-architecture - Infrastructure layer setup