在微服务架构里,业务数据落库和向消息中间件发送事件常常分属不同资源,传统做法先提交数据库再发消息,一旦进程崩溃就会出现数据已写但消息未发,或者消息发出但回滚了的尴尬局面。Outbox发件箱模式把事件当作普通行写入同一数据库的专用表,利用本地事务保证业务与消息记录同时成功或失败,随后由另一个进程读取发件箱并投递到MQ,从而实现分布式下的准实时一致。

一、Outbox模式的核心原理
Outbox的本质是“先存后发”。在单次数据库事务中,除了更新订单、账户等聚合根,还向Outbox表插入一条包含消息体、目标主题和唯一消息编号的记录。因为都在同一个本地事务里,要么全成功要么全失败,从根源上消灭了跨资源的不原子问题。事务提交后,业务接口即可返回,不需要等待消息真正到达Broker。
随后启动的投递服务(Delivery Worker)定期扫描Outbox表中状态为“待发送”的记录,将其推送到RabbitMQ或Kafka,发送成功后再把该行标记为“已发送”或物理删除。即便投递器中途挂掉,下次启动仍能从表里捡回未完成的消息,因此消息至少会被投递一次。配合消费者侧幂等,就能做到不丢不重。
1.1 表结构设计
Outbox表至少需要消息ID、聚合类型、 payload、创建时间和状态字段。下面给出SQL Server的建表语句,消息ID使用UUID避免多实例冲突,payload用NVARCHAR(MAX)存放序列化后的JSON。
CREATE TABLE OutboxMessage (
Id UNIQUEIDENTIFIER PRIMARY KEY,
AggregateType NVARCHAR(100) NOT NULL,
Topic NVARCHAR(200) NOT NULL,
Payload NVARCHAR(MAX) NOT NULL,
CreatedAt DATETIME2 NOT NULL DEFAULT SYSUTCDATETIME(),
Status INT NOT NULL DEFAULT 0 -- 0:Pending, 1:Sent, 2:Failed
);
Status字段便于做可视化监控和失败重试,如果希望减少表体积,也可以在发送成功后直接删除行,但保留少量历史有助于排查问题。索引上建议给Status和CreatedAt建组合索引,提升轮询效率。
二、C#中事务内双写实现
使用ADO.NET或Dapper都可以在一个SqlTransaction里完成业务与Outbox的插入。下面示例采用Dapper,在创建订单的同时写入Outbox,最后统一提交。
public class OrderService
{
private readonly string _connStr;
public OrderService(string connStr) { _connStr = connStr; }
public void CreateOrderWithEvent(Order order)
{
using var conn = new SqlConnection(_connStr);
conn.Open();
using var tx = conn.BeginTransaction();
try
{
// 1. 写业务表
conn.Execute(
"INSERT INTO Orders(Id,Amount,Status) VALUES(@Id,@Amount,@Status)",
order, transaction: tx);
// 2. 写Outbox表
var msg = new OutboxMessage
{
Id = Guid.NewGuid(),
AggregateType = "Order",
Topic = "order.created",
Payload = JsonSerializer.Serialize(order)
};
conn.Execute(
"INSERT INTO OutboxMessage(Id,AggregateType,Topic,Payload) VALUES(@Id,@AggregateType,@Topic,@Payload)",
msg, transaction: tx);
tx.Commit();
}
catch
{
tx.Rollback();
throw;
}
}
}
这段代码的关键在于两个INSERT共用同一个tx,任何一步异常都会整体回滚。相比先提交订单再发MQ,这里接口耗时只多了一次本地表插入,性能影响极小,却换来了强一致保障。
如果你使用Entity Framework Core,可以把Outbox实体加入DbContext,在SaveChanges前追加消息实体,利用SaveChanges的一次事务提交达成同样效果。无论哪种ORM,核心原则都是“同一事务、同一连接”。
2.1 消息体版本与序列化
Payload建议带上Schema版本号,例如{ "v": 1, "data": { ... } },方便消费者在字段变更时兼容旧消息。C#端可用record类型配合System.Text.Json,避免二进制格式化带来的耦合。
三、独立的消息投递器
投递器是一个长运行后台任务,可以用IHostedService实现,每隔几秒拉取待发送消息并推送到RabbitMQ。下面代码展示基础轮询逻辑。
public class OutboxDispatcher : BackgroundService
{
private readonly string _connStr;
private readonly IModel _channel;
public OutboxDispatcher(string connStr, IModel channel)
{
_connStr = connStr;
_channel = channel;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
using var conn = new SqlConnection(_connStr);
var pending = (await conn.QueryAsync<OutboxMessage>(
"SELECT TOP 50 * FROM OutboxMessage WHERE Status=0 ORDER BY CreatedAt"))
.ToList();
foreach (var m in pending)
{
var body = Encoding.UTF8.GetBytes(m.Payload);
_channel.BasicPublish("", m.Topic, null, body);
await conn.ExecuteAsync(
"UPDATE OutboxMessage SET Status=1 WHERE Id=@Id",
new { m.Id });
}
await Task.Delay(2000, stoppingToken);
}
}
}
上述实现保证了消息至少投递一次。若RabbitMQ不可用,更新Status会失败,下次循环继续重试。生产环境应增加失败计数和退避策略,避免对故障Broker狂打日志。
为了提升吞吐,可把拉取和发送改为批量发布,或在多实例部署时使用抢占式锁(如UPDATE TOP 50 WITH (UPDLOCK, READPAST))防止重复发送。消费者端则依据消息Id做幂等表去重,彻底消除重复副作用。
3.1 与Transactional Outbox框架对比
社区有NServiceBus、MassTransit等库内置了Outbox,它们自动拦截命令并管理发件箱,减少手写代码。但自研实现更轻量,也更容易贴合已有仓储层。选择时权衡团队运维能力与定制需求即可。
四、常见陷阱与优化建议
第一个陷阱是投递器与业务写表争用锁。高频写入场景下,扫描Outbox会触发表扫描,应在Status和CreatedAt上建索引,并控制每次TOP数量。第二个陷阱是消息膨胀,长期不清理已发送记录会让表变大,可定时归档或改为发送后即删。
另一个重点是消费者幂等。因为Outbox只保证至少一次,网络重发或投递器重启都可能重复,消费方必须用消息Id或业务键做唯一约束。例如订单积分服务在发放前先INSERT IGNORE一条消费记录,确保同一事件只生效一次。
| 方案 | 一致性 | 复杂度 | 适用场景 |
|---|---|---|---|
| 本地事务表+轮询 | 最终一致 | 低 | 中小规模服务 |
| MQ事务消息 | 最终一致 | 中 | Broker支持事务 |
| 2PC | 强一致 | 高 | 金融核心链路 |
Outbox以较低复杂度换取了接近无忧的可靠性,是C#后端处理跨服务通知的首选。落地时只要守住“同事务写、独立发、消费幂等”三条线,就能平稳替代易错的双写旧逻辑。