From d5de65fdb1802cdf498d65d93397f290813c377c Mon Sep 17 00:00:00 2001 From: muqiuhan Date: Tue, 9 Sep 2025 06:17:00 +0000 Subject: deploy: 9665097f0fa0dae9f92123fac54d26c0818758a5 --- 2025/05/07/nestjs-bullmq-mail-business/index.html | 41 ++++++----- .../index.html | 82 ++++++++++++++-------- 2025/05/27/vertical-slicing-practice/index.html | 2 +- 3 files changed, 75 insertions(+), 50 deletions(-) (limited to '2025/05') diff --git a/2025/05/07/nestjs-bullmq-mail-business/index.html b/2025/05/07/nestjs-bullmq-mail-business/index.html index 6b36afbd..5eb3bcbd 100644 --- a/2025/05/07/nestjs-bullmq-mail-business/index.html +++ b/2025/05/07/nestjs-bullmq-mail-business/index.html @@ -192,48 +192,53 @@
-

在 BullMQ(以及它在 NestJS 里包装的 @Processor/WorkerHost)里,整个生命周期大致是这样的:

+

在 BullMQ(以及它在 NestJS 里包装的 @Processor/WorkerHost)里,整个生命周期大致是这样的:

    -
  1. 队列(在 NestJS 里由 @Processor 装饰的类)会被一个底层的 Worker 订阅。
  2. -
  3. 有新任务(job)进来时,Worker 会调用写在该类里的 async process(job: Job) 方法。
  4. -
  5. 如果 process() 正常返回(即没有抛异常),Job 就被标记为 completed,然后才会去触发所有注册了 @OnWorkerEvent('completed') 的回调。
  6. +
  7. 队列(在 NestJS 里由 @Processor 装饰的类)会被一个底层的 Worker 订阅。
  8. +
  9. 有新任务(job)进来时,Worker 会调用写在该类里的 async process(job: Job) 方法。
  10. +
  11. 如果 process() 正常返回(即没有抛异常),Job 就被标记为 completed,然后才会去触发所有注册了 @OnWorkerEvent('completed') 的回调。

也就是说:

而我在此处的业务目的是 “用队列来做可靠的、可重试的邮件发送”,那么一定要把发送邮件的逻辑写到 process() 里,这样在 commandBus.execute(new SendMailCommand(...)) 抛错时,BullMQ 会根据创建 JOB 时的重试策略(retry、backoff 等)自动重新入队。而把它放到 onCompleted(),只相当于 job 成功完成后的“事后通知”,一旦失败不会再重试,也无法利用 BullMQ 的锁、超时、重试机制。

-

举个最简化的调整示例,删掉 onCompleted,把真正的发信放到 process:

+

举个最简化的调整示例,删掉 onCompleted,把真正的发信放到 process:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// ... existing imports ...

@Processor(process.env.MAILER_QUEUE_NAME || "gcpm-mailer")
export class BullMQMailerProcesser extends WorkerHost {
constructor(
private readonly commandBus: CommandBus,
private readonly logger: LoggingService,
) {
super();
}

// ① 当有新 job 拉取到时,这个方法会被调用
public async process(job: Job): Promise<void> {
const mailAggregate = new Mail(job.data.mail);
try {
await this.commandBus.execute(new SendMailCommand(mailAggregate));
} catch (err) {
this.logger.error(`邮件发送失败,jobId=${job.id}`, err);
// 抛出错误,触发重试或失败
throw err;
}
}

// ② onCompleted 仅在 process() 正常返回后触发,
// 不建议在这里执行核心业务(也无法触发重试)。
// @OnWorkerEvent("completed")
// async onCompleted(job: Job) { … }
}
-

参考 BullMQ 官方文档:


而 重试次数本身并没有一个硬性上限,完全由添加 Job 时通过 attempts 这个选项来控制:

示例(给某封邮件最多重试 3 次):

1
2
3
4
5
6
7
8
9
10
11
await this.mailerQueue.add(
id,
{ mail: new Mail(/*…*/ ) },
{
attempts: 3, // 最多尝试 3 次
backoff: { // 重试时的延迟策略(可选)
type: 'exponential',
delay: 1000,
},
},
);
-

还有一个需要注意的地方,在我的业务中,邮件发送的是一种时间区间报告,这个报告包含了过去二十四小时的一些系统中的事件,但如果重试有延迟策略或重试本身就有计算成本的话,这封邮件就不是 “过去二十四小时” 的了,因为重试带来了一个真空期。

换言之,这个问题本质上是——重试导致「发送时刻」与「原始 24 小时窗口」错开,从而让邮件里报出来的数据不再精确。常见的解决思路就是:把「窗口定义」或者「报表内容」在调度时就固化下来,真正的队列任务只负责发送,而不再实时去重新计算时间区间。

我想到了两种解决方案:

-

一、任务参数里带上「时间区间」
在 enqueue 的时候,就算出 windowStart/windowEnd,然后把它放到 job.data 里。无论后面 process 什么时候真正跑,都是基于同一个时间区间去查询:

+

一、任务参数里带上「时间区间」
+在 enqueue 的时候,就算出 windowStart/windowEnd,然后把它放到 job.data 里。无论后面 process 什么时候真正跑,都是基于同一个时间区间去查询:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
// 调度时
const now = new Date();
const windowStart = new Date(now.getTime() - 24 * 60 * 60 * 1000);
await this.mailerQueue.add(
id,
{
mail: new Mail({
...options,
id,
sentAt: now,
status: MailStatus.PENDING,
windowStart,
windowEnd: now,
}),
},
{
attempts: 3,
backoff: { type: 'exponential', delay: 1000 },
},
);

// process 里
public async process(job: Job) {
const { windowStart, windowEnd } = job.data.mail;
// ① 只查询 [windowStart, windowEnd] 的事件
const events = await this.reportService.findEvents(windowStart, windowEnd);
const reportHtml = await this.reportService.renderReport(events);
await this.commandBus.execute(new SendMailCommand(job.data.mail, reportHtml));
}
-

➜ 这样无是马上执行还是几次重试后才执行,数据规则都不会变。

-

二、预先生成「静态报表内容」,挂到队列里
如果计算成本很高,或者怕重复查询数据开销大,也可以在调度时就把最终的 HTML/Text/附件 都先打好,然后作为 job.data 传进去,真正的 process() 只做一次“发送”即可:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// 调度时:先生成报告
const now = new Date();
const windowStart = new Date(now.getTime() - 24*3600*1000);
const events = await this.reportService.findEvents(windowStart, now);
const reportHtml = await this.reportService.renderReport(events);

// 把静态内容塞到队列
await this.mailerQueue.add(
id,
{
mail: new Mail({ /*…*/, windowStart, windowEnd: now }),
reportHtml, // <- 预渲染好的文本/HTML
attachments: […], // <- 如果有附件也一并塞
},
{ attempts: 3, backoff: { type: 'fixed', delay: 5_000 } },
);

// process 里只关注发送
public async process(job: Job) {
try {
await this.mailService.send({
to: job.data.mail.to,
subject: `系统 24h 报表`,
html: job.data.reportHtml,
attachments: job.data.attachments,
});
} catch (e) {
throw e; // 触发重试
}
}

➜ 重试带来的任何延迟,都不影响邮件正文,始终是一份「事先约定好、并且静态化」的报告。

-

这两种模式都能保证最终发送时的数据窗口或内容,与当初调度时的预期完全一致,不会因为重试延迟而出现“数据真空”或“多算/少算”问题。

-

参考文档:

diff --git a/2025/05/08/Multiplayer-Collaborative-Systems-tips/index.html b/2025/05/08/Multiplayer-Collaborative-Systems-tips/index.html index 47a9feef..42f2895c 100644 --- a/2025/05/08/Multiplayer-Collaborative-Systems-tips/index.html +++ b/2025/05/08/Multiplayer-Collaborative-Systems-tips/index.html @@ -193,36 +193,41 @@

最近碰到一块业务:在系统中可以存在多个用户同时对某个项目信息进行编辑,这种多人协作的场景挺有意思的,不过在我们的业务中,并不需要实时协作,只需要保证不会出错就行,话虽如此,但也可以探索一下实时协作的实现方案,防止老年痴呆。

-

先来看看第一个方案 —— CRDT(Conflict-free Replicated Data Type,无冲突可复制数据类型)是一类数据结构,它保证了在分布式节点(或多客户端)上进行离线/并发更新后,无需中心协调、也无需人工干预,通过“合并策略”就能得到一致的最终状态。

+

先来看看第一个方案 —— CRDT(Conflict-free Replicated Data Type,无冲突可复制数据类型)是一类数据结构,它保证了在分布式节点(或多客户端)上进行离线/并发更新后,无需中心协调、也无需人工干预,通过“合并策略”就能得到一致的最终状态。

核心思想是:所有并发操作都是幂等(idempotent)、可交换(commutative)的。

-

常见类型有:
一、G-Counter(只能增计数器)
二、PN-Counter(可增可减计数器)
三、LWW-Register(最后写入胜出)
四、结合 JSON 的树型 CRDT(如 Automerge / Yjs)

-

更多原理可参考 Decipad 博客“Collaborative and Offline Editing Using CRDTs”[^1]。

+

常见类型有:
+一、G-Counter(只能增计数器)
+二、PN-Counter(可增可减计数器)
+三、LWW-Register(最后写入胜出)
+四、结合 JSON 的树型 CRDT(如 Automerge / Yjs)

+

更多原理可参考 Decipad 博客“Collaborative and Offline Editing Using CRDTs”[1]。

有一个挺有趣的 Rust 项目 Loro: Make your JSON data collaborative and version-controlled with CRDTs

-

假设我的项目信息编辑页面允许多人实时/离线修改某个研究项目的“名称”、“描述”字段,前端用 SvelteKit + GraphQL 获取和提交变更:

+

假设我的项目信息编辑页面允许多人实时/离线修改某个研究项目的“名称”、“描述”字段,前端用 SvelteKit + GraphQL 获取和提交变更:

1
2
┌── 用户 A 离线修改了 “description” 的若干段文本  
└── 用户 B 同时在线修改了同一字段的其他段落
-

如果后端使用 CRDT(比如把 description 用 JSON-CRDT 存储),两次修改只要在任意顺序合并都能得到完整的内容:

-

首先,A 客户端本地 apply 操作并缓存,恢复网络后推给服务器;
然后,服务器用 CRDT merge(A.delta, B.delta),得到一致文档
最后,服务器广播新文档到所有客户端,A/B 均得到相同结果

+

首先,A 客户端本地 apply 操作并缓存,恢复网络后推给服务器;
+然后,服务器用 CRDT merge(A.delta, B.delta),得到一致文档
+最后,服务器广播新文档到所有客户端,A/B 均得到相同结果


-

好了说点实际符合业务场景的方案,首先想到的是悲观锁(Pessimistic Locking) ,思路是:用户打开编辑界面时,向后端申请“锁” → 其它用户尝试编辑时被拒绝 → 编辑完成后释放锁/超时自动释放。

+

好了说点实际符合业务场景的方案,首先想到的是悲观锁(Pessimistic Locking) ,思路是:用户打开编辑界面时,向后端申请“锁” → 其它用户尝试编辑时被拒绝 → 编辑完成后释放锁/超时自动释放。

假如有一个这样的锁表:

-
CREATE TABLE project_lock (
+
CREATE TABLE project_lock (
     project_id UUID PRIMARY KEY,
     locked_by  UUID NOT NULL,
     expires_at TIMESTAMPTZ NOT NULL
 );
 

可以在事务内申请它:

-
const now = new Date();
+
const now = new Date();
 const expires = new Date(now.getTime() + 5*60*1000); // 5 分钟后过期
 await prisma.$transaction(async tx => {
     const existing = await tx.project_lock.findUnique({ where:{ project_id } });
         
     if (existing && existing.expires_at > now) {
-        throw new Error('项目正被人编辑');
-    }
+	    throw new Error('项目正被人编辑');
+	}
         
     await tx.project_lock.upsert({
         where: { project_id },
@@ -232,14 +237,16 @@ await prisma.$transaction(async tx => {
 });
 

释放锁就直接从锁表里删掉对应的数据即可:

-
await prisma.project_lock.delete({ where:{ project_id } });
+
await prisma.project_lock.delete({ where:{ project_id } });
 

前端的话,大概就是:

-

在进入编辑前请求一下 /api/project/:id/lock 之类的 API,失败则提示“被占用”;
在 onbeforeunload 时执行 /unlock;
超时后后端自动允许新锁。

+

在进入编辑前请求一下 /api/project/:id/lock 之类的 API,失败则提示“被占用”;
+在 onbeforeunload 时执行 /unlock;
+超时后后端自动允许新锁。


-

第二个方案是乐观并发控制(Optimistic Concurrency) :记录资源的版本号或时间戳;客户端提交更新时带上自己的版本号,后端检查版本是否一致,不一致则认为冲突,返回 409,由客户端告知用户“数据已过期,请刷新后合并”。

+

第二个方案是乐观并发控制(Optimistic Concurrency) :记录资源的版本号或时间戳;客户端提交更新时带上自己的版本号,后端检查版本是否一致,不一致则认为冲突,返回 409,由客户端告知用户“数据已过期,请刷新后合并”。

具体实现中,可以尝试在 project 表加上 version INT NOT NULL DEFAULT 1, updated_at TIMESTAMPTZ ,然后更新项目时:

-
async updateProject(parent, { id, version, input }, ctx) {
+
async updateProject(parent, { id, version, input }, ctx) {
     const result = await prisma.$executeRaw`
     UPDATE project
         SET name        = ${input.name},
@@ -247,38 +254,51 @@ await prisma.$transaction(async tx => {
             version     = version + 1,
             updated_at  = now()
     WHERE id = ${id} AND version = ${version}
-    `;
+	`;
     if (result === 0) {
-        throw new ConflictException('版本冲突,请刷新后重试');
+	    throw new ConflictException('版本冲突,请刷新后重试');
     }
     return prisma.project.findUnique({ where:{ id } });
 }
 
-

前端捕获到冲突错误可以用一个弹窗提示“另有用户已更新此项目,是否合并/重新加载?” 之类的玩意儿。

+

前端捕获到冲突错误可以用一个弹窗提示“另有用户已更新此项目,是否合并/重新加载?” 之类的玩意儿。


-

第三个方案是:操作转化(Operational Transformation,OT)

-

也就是记录用户每次的“操作”(insert/delete at position),服务器根据历史操作序列对并发操作做转化(transform),确保先到达的操作调整后再应用后到达的。

-

有一些实现案例:
一、ShareDB(Node.js)
二、Google Docs 中的同步算法

-

具体实现的话,可能要现在前端逐字符/块地包装成操作并 WebSocket 推送,服务器再维护一个“操作历史队列”,每来一个 op 就 transform 并 broadcast,而客户端收到广播后,按顺序 replay 保证视图一致。

+

第三个方案是:操作转化(Operational Transformation,OT)

+

也就是记录用户每次的“操作”(insert/delete at position),服务器根据历史操作序列对并发操作做转化(transform),确保先到达的操作调整后再应用后到达的。

+

有一些实现案例:
+一、ShareDB(Node.js)
+二、Google Docs 中的同步算法

+

具体实现的话,可能要现在前端逐字符/块地包装成操作并 WebSocket 推送,服务器再维护一个“操作历史队列”,每来一个 op 就 transform 并 broadcast,而客户端收到广播后,按顺序 replay 保证视图一致。


最后可能还可以用事件溯源(Event Sourcing)+ 场景命令模式来实现:

-

不直接存状态,而是存所有“命令 / 事件”(Event),回放事件得到当前状态。冲突通过合并策略或补偿事件(Compensating Events)解决。

-

例如:
在每次更新时推送 ProjectUpdated { projectId, fieldsChanged, userId, timestamp } ,
然后写入事件存储(如 Kafka / EventStoreDB),
读端 Consumer 按顺序重建最新状态或按领域聚合 ,
最后在并发时如果两个事件都修改了同一字段,可在写端做校验/补偿,或在读端做最后写入胜出等策略 。

+

不直接存状态,而是存所有“命令 / 事件”(Event),回放事件得到当前状态。冲突通过合并策略或补偿事件(Compensating Events)解决。

+

例如:
+在每次更新时推送 ProjectUpdated { projectId, fieldsChanged, userId, timestamp } ,
+然后写入事件存储(如 Kafka / EventStoreDB),
+读端 Consumer 按顺序重建最新状态或按领域聚合 ,
+最后在并发时如果两个事件都修改了同一字段,可在写端做校验/补偿,或在读端做最后写入胜出等策略 。


总结来说,

    -
  • CRDT 最擅长 去中心化、离线编辑、自动合并;
  • -
  • 若不引入 CRDT,可根据业务侧重点选用:
      -
    1. 悲观锁 → 强制串行编辑,简单粗暴;
    2. -
    3. 乐观并发 → 适合大多数业务场景,成本低;
    4. -
    5. OT → 适合富文本或实时协同场景,复杂度中等;
    6. +
    7. CRDT 最擅长 去中心化、离线编辑、自动合并;
    8. +
    9. 若不引入 CRDT,可根据业务侧重点选用: +
        +
      1. 悲观锁 → 强制串行编辑,简单粗暴;
      2. +
      3. 乐观并发 → 适合大多数业务场景,成本低;
      4. +
      5. OT → 适合富文本或实时协同场景,复杂度中等;
      6. 事件溯源 → 适合需要全历史审计、可回放的场景。

-

[^1]: Decipad 博客 “Collaborative and Offline Editing Using CRDTs”
https://www.decipad.com/blog/decipads-innovative-method-collaborative-and-offline-editing-using-crdts

-

[^2]: Hacker News 讨论(CRDT 相关线程)
https://news.ycombinator.com/item?id=38289327

+
+
+
    +
  1. Decipad 博客 “Collaborative and Offline Editing Using CRDTs”
    +https://www.decipad.com/blog/decipads-innovative-method-collaborative-and-offline-editing-using-crdts ↩︎

    +
  2. +
+
diff --git a/2025/05/27/vertical-slicing-practice/index.html b/2025/05/27/vertical-slicing-practice/index.html index 98dee1bb..206a3e98 100644 --- a/2025/05/27/vertical-slicing-practice/index.html +++ b/2025/05/27/vertical-slicing-practice/index.html @@ -210,7 +210,7 @@

对于登录功能,需要考虑:

  • UI 层:登录表单(输入邮箱、密码的地方)、提交按钮、错误提示信息。
  • -
  • API/服务层:接收登录请求、验证用户凭证的接口。
  • +
  • API/服务层:接收登录请求、验证用户凭证的接口。
  • 业务逻辑层:校验输入格式、查询用户信息、验证密码、生成会话(Session)或令牌(Token)。
  • 数据访问层:从数据库中读取用户信息。
-- cgit v1.2.3