MikroORM 事务性发件箱模式实战:将领域事件与业务数据写入同一事务,实现可靠的事件发布

原创2026-09-26 12:19:011,173 阅读
文章标签:后端

MikroORM 事务性发件箱模式实战:将领域事件与业务数据写入同一事务,实现可靠的事件发布

构建事件驱动应用时,一个经典难题是如何保证“业务数据落库”与“领域事件发布”的一致性。本文基于 MikroORM 官方指南《Transactional Outbox Pattern》,完整讲解事务性发件箱(Transactional Outbox)模式的落地方式:从 outbox 实体的四种定义写法、事务内写入事件、用 onFlush 订阅者自动收集变更集,到发布 worker 的轮询、幂等与悲观锁处理。读完本文,你将掌握一套可在 MikroORM 项目中直接复制运行的事件可靠发布方案。

为什么需要事务性发件箱:事件丢失窗口

在事件驱动架构中,我们通常希望在 UserCreated、OrderPlaced 这类领域事件发生时,把它们发布到消息中间件(如 Kafka、RabbitMQ 或事件总线)。一个很自然的想法是:在数据库事务提交成功后触发钩子,再发布事件。例如在 afterTransactionCommit 钩子中发布。

但这种“先提交、后发布”的朴素做法存在一个危险的时间窗口:如果进程在事务提交之后、事件发布之前崩溃,事件就永久丢失了,因为没有任何东西被持久化下来,事后无法恢复,最终导致业务状态与外部系统不一致。

事务性发件箱模式通过改变持久化边界来消除这个窗口:

  1. 在业务事务中,除了写入业务数据,把待发布的事件也写入一张 outbox 表;
  2. 业务数据与事件行共享同一个事务——要么一起提交,要么一起回滚;
  3. 由另一个独立进程(cron、后台 worker 或独立微服务)轮询 outbox 表,把事件发布到消息中间件,并标记为已处理。

这样一来,事件在事务提交时就已经落库,进程崩溃也不会丢失事件,系统获得至少一次(at-least-once)投递保证。代价是消费者必须容忍重复消息,即要求幂等消费。

定义 OutboxEvent 实体

outbox 事件本身就是一个普通实体。MikroORM 支持多种实体定义方式,以下四种写法在功能上等价,你可以根据项目使用的元数据提供器选择其一。四段代码均放在 ./entities/OutboxEvent.ts 中。

方式一:defineEntity + class(推荐)

import { defineEntity, p } from '@mikro-orm/core';

const OutboxEventSchema = defineEntity({
  name: 'OutboxEvent',
  properties: {
    id: p.integer().primary(),
    eventType: p.string(),
    payload: p.json<Record<string, unknown>>(),
    createdAt: p.datetime().onCreate(() => new Date()),
    processed: p.boolean().default(false),
  },
});

export class OutboxEvent extends OutboxEventSchema.class {}
OutboxEventSchema.setClass(OutboxEvent);

这种写法把 Schema 定义与 class 声明解耦:先通过 defineEntity 声明元数据,再创建一个空 class 并回填给 Schema。它不依赖 reflect-metadata,类型推断依然完整。

方式二:纯 defineEntity(无 class)

import { type InferEntity, defineEntity, p } from '@mikro-orm/core';

export const OutboxEvent = defineEntity({
  name: 'OutboxEvent',
  properties: {
    id: p.integer().primary(),
    eventType: p.string(),
    payload: p.json<Record<string, unknown>>(),
    createdAt: p.datetime().onCreate(() => new Date()),
    processed: p.boolean().default(false),
  },
});

export type IOutboxEvent = InferEntity<typeof OutboxEvent>;

不需要 class 时,可以直接把 Schema 作为实体导出,并用 InferEntity 推导出对应的实体类型,适合纯数据驱动的代码风格。

方式三:reflect-metadata 装饰器

import { Entity, PrimaryKey, Property } from '@mikro-orm/core';

@Entity()
export class OutboxEvent {

  @PrimaryKey()
  id!: number;

  @Property()
  eventType!: string;

  @Property({ type: 'json' })
  payload!: Record<string, unknown>;

  @Property()
  createdAt = new Date();

  @Property()
  processed = false;

}

经典装饰器写法,依赖 reflect-metadata 提供类型反射,@Property({ type: 'json' }) 显式指定 JSON 列类型以保存任意结构化载荷。

方式四:ts-morph(代码生成)

import { Entity, PrimaryKey, Property } from '@mikro-orm/core';

@Entity()
export class OutboxEvent {

  @PrimaryKey()
  id!: number;

  @Property()
  eventType!: string;

  @Property()
  payload!: Record<string, unknown>;

  @Property()
  createdAt = new Date();

  @Property()
  processed = false;

}

ts-morph 元数据提供器通常在 @mikro-orm/reflection 的代码生成场景下使用,写法与装饰器版一致,由工具在编译期提取元数据,无需运行时反射。

四种方式的关键字段语义一致:

字段 类型 说明
id integer + primary 主键,自增整数
eventType string 事件类型标识,如 User_create
payload json 事件载荷,任意 JSON 对象
createdAt datetime 创建时间,onCreate 在插入时自动填充当前时间
processed boolean 是否已发布,默认 false

在事务内写入事件

由于 outbox 事件就是普通实体,你可以把它和业务实体一起持久化。它们会参与同一次 em.flush() 事务:

await em.transactional(async em => {
  const user = em.create(User, { name, email });

  em.create(OutboxEvent, {
    eventType: 'User_create',
    payload: { name, email },
  });
});
// 此时 user 行与 outbox 行要么同时提交,要么(事务失败时)同时回滚

这段代码的核心保证在于:em.transactional 的回调结束时会自动触发 flush,把回调内 em.create 产生的所有待持久化实体一次性写入数据库。从源码看,transactional 方法(EntityManager.ts)在回调结束时统一提交,从而让 User 与 OutboxEvent 落在同一个事务边界内。这正是 outbox 模式“事件与业务数据同生共死”的关键。

自动化:用 onFlush 订阅者收集变更集

如果每个事务都手动 em.create(OutboxEvent, ...),代码会非常啰嗦且容易遗漏。更优雅的做法是注册一个 onFlush 订阅者,在 flush 过程中检查工作单元(Unit of Work)计算出的变更集(ChangeSet),为每一条业务变更自动排队一条 outbox 事件:

import { ChangeSetType, EventSubscriber, FlushEventArgs } from '@mikro-orm/core';
import { OutboxEvent } from './entities/OutboxEvent.js';

export class OutboxSubscriber implements EventSubscriber {

  async onFlush(args: FlushEventArgs) {
    for (const cs of args.uow.getChangeSets()) {
      if (cs.meta.className === 'OutboxEvent') {
        continue; // 不要为 outbox 事件本身再创建 outbox 事件
      }

      // 跳过用于自引用关系处理的内部早更新变更集
      if (cs.type === ChangeSetType.UPDATE_EARLY) {
        continue;
      }

      const event = args.em.create(OutboxEvent, {
        eventType: `${cs.meta.className}_${cs.type}`,
        payload: cs.getPrimaryKey(true),
      });

      args.uow.computeChangeSet(event);
    }
  }

}

然后在 ORM 配置中注册订阅者:

MikroORM.init({
  subscribers: [new OutboxSubscriber()],
});

订阅者背后的 flush 时序

理解这段代码,需要知道 onFlush 在 flush 生命周期中的位置。查看 UnitOfWork.ts 中 doCommit() 的实现,可以发现 flush 的完整顺序是:

  1. 触发 beforeFlush 事件;
  2. computeChangeSets() —— 计算实体当前状态与持久化状态的差异,生成变更集;
  3. 触发 onFlush 事件(UnitOfWork.ts)—— 订阅者在这里读取已计算的变更集;
  4. 处理集合更新;
  5. persistToDatabase() 把变更集分组写入数据库(若当前不在事务中且开启了隐式事务,会用 connection.transactional 包装整个持久化过程);
  6. 触发 afterFlush 事件。

因此,OutboxSubscriber.onFlush 中通过 args.uow.getChangeSets() 拿到的,正是本次 flush 即将落库的全部实体变更。订阅者遍历这些变更集,为每条业务变更生成一条 OutboxEvent,再调用 args.uow.computeChangeSet(event) 手动把新实体纳入工作单元的变更集计算,从而让这条 outbox 事件与业务变更一起进入后续的 persistToDatabase 阶段,在同一事务中落库。

代码里有两个值得注意的防御性判断:

  • 跳过 OutboxEvent 自身:否则每条 outbox 事件又会产生一条新的 outbox 事件,形成无限递归;
  • 跳过 ChangeSetType.UPDATE_EARLY:这是框架内部用于处理自引用关系的早更新变更集,不属于业务变更,不应产生领域事件。

EventSubscriber 接口与 FlushEventArgs 的定义位于 EventSubscriber.ts,onFlush 支持同步或异步实现。

发布事件:独立 worker 轮询 outbox 表

事件写入 outbox 表后,需要一个独立进程把事件发出去。它可以是一个 cron 任务、后台 worker,也可以是专门的微服务:

async function processOutbox(orm: MikroORM) {
  const em = orm.em.fork();

  const events = await em.find(
    OutboxEvent,
    { processed: false },
    { orderBy: { createdAt: 'ASC' }, limit: 100 },
  );

  for (const event of events) {
    await publishToMessageBroker(event.eventType, event.payload);
    event.processed = true;
  }

  await em.flush();
}

要点说明:

  • orm.em.fork() 创建一个独立上下文的 EntityManager,避免污染全局上下文,也避免长生命周期上下文中的身份映射(Identity Map)问题;
  • 按 createdAt 升序、limit: 100 分批拉取未处理事件,避免一次处理过多导致内存压力;
  • 每条事件发布成功后置 processed = true,最后统一 flush() 一次提交。

消费者必须幂等

由于 outbox 模式提供的是至少一次投递而非精确一次投递,同一事件可能被重复发布。典型场景:worker 发布了部分事件后、在 flush() 标记 processed 之前崩溃,下一轮运行会重新发布这些事件。因此消费端必须设计为幂等——例如以事件 ID 或业务主键做去重,确保重复消息不产生副作用。

多 worker 并发:使用悲观锁跳过已锁行

如果部署了多个发布实例,两个 worker 可能同时读到同一批未处理事件并重复发布。此时应使用数据库悲观锁,让并发查询跳过已被其他事务锁定的行,即 SQL 的 FOR UPDATE SKIP LOCKED 语义。在 MikroORM 中对应 LockMode.PESSIMISTIC_PARTIAL_WRITE:

import { LockMode } from '@mikro-orm/core';

const events = await em.find(
  OutboxEvent,
  { processed: false },
  {
    orderBy: { createdAt: 'ASC' },
    limit: 100,
    lockMode: LockMode.PESSIMISTIC_PARTIAL_WRITE,
  },
);

从 enums.ts 中可以看到完整的锁模式枚举,PESSIMISTIC_PARTIAL_WRITE 被定义为排他锁且跳过已锁行(FOR UPDATE SKIP LOCKED),是并发轮询场景的标准选择。相关的锁模式还包括:

锁模式 SQL 语义 适用场景
PESSIMISTIC_READ FOR SHARE 共享读锁
PESSIMISTIC_WRITE FOR UPDATE 排他写锁,阻塞等待
PESSIMISTIC_PARTIAL_WRITE FOR UPDATE SKIP LOCKED 跳过已锁行,适合并发 worker 轮询
PESSIMISTIC_WRITE_OR_FAIL FOR UPDATE NOWAIT 遇到锁立即失败

注意:SKIP LOCKED 需要数据库支持(PostgreSQL、MySQL 8+、SQLite 3.8.3+ 等),且通常在显式事务或可锁定查询上下文中生效,实际行为以所用方言为准。

清理已处理事件

outbox 表会持续增长,需要定期清理。用 nativeDelete 直接删除 7 天前已处理的事件:

const cutoff = new Date();
cutoff.setDate(cutoff.getDate() - 7);

await em.nativeDelete(OutboxEvent, {
  processed: true,
  createdAt: { $lt: cutoff },
});

$lt 是 MikroORM 查询条件操作符,表示“小于”;nativeDelete 绕过实体生命周期钩子直接执行删除,适合批量清理场景。注意:保留窗口(此处为 7 天)应根据下游消费与重放需求调整,并建议为该查询建立 (processed, createdAt) 联合索引。

为什么不直接在钩子里发布事件

一个更“省事”的诱惑是把事件发布写进 afterFlush 或 afterTransactionCommit 钩子。这种方案代码更少,但有致命缺陷:如果进程在事务提交之后、事件发出之前崩溃,事件就彻底丢失了,而且由于没有持久化任何记录,事后没有任何恢复手段。

outbox 模式通过把事件持久化纳入事务本身来消除这个窗口——事件要么随业务数据一起提交,要么一起回滚,进程崩溃最坏只会导致事件被重复发布,而这可以通过幂等消费者轻松消化。这正是事务性发件箱模式在可靠事件发布场景中被广泛采用的根本原因。

源码与文档定位

小结

事务性发件箱模式在 MikroORM 中的落地可以总结为三步:同一事务写事件(用 em.transactional 或 onFlush 订阅者自动生成 outbox 记录)、独立进程发布事件(轮询 + 标记 processed)、工程化加固(幂等消费者、FOR UPDATE SKIP LOCKED 并发锁、定期清理)。它把“事件发布”从脆弱的内存操作变成可恢复的持久化流程,是构建高可靠事件驱动系统的实用基础。

登录后查看全文
mikro-orm