后台任务、消息与可靠流程
连接持久任务、事务发件箱、消费者幂等、调度、重试与服务消息边界。
API 接受后由谁继续
请求可以同步等结果,也可以返回已接受的任务身份。后台消费需要长驻 worker 或具明确完成保证的托管执行环境。普通 serverless 请求返回后,进程内 Promise 可能被暂停或结束;平台 waitUntil 等接口有自己的期限与范围,不能代替永久队列。
BullMQ 将任务与状态放到相应后端,生产者投递,worker 执行。Redis 路线需要连接、持久化、内存与逐出策略;PostgreSQL 队列如 pg-boss、Graphile Worker 提供另一种基础设施组合。采用时读取该版本的调度和数据模式,不绕过库接口直接修改内部表。
业务提交与入队之间的间隙
先提交业务表,再调用队列,会在二者之间留下进程崩溃或网络失败窗口。业务已经改变,任务却没记录;先入队再提交则可能让任务执行尚未成立的业务。两份存储没有共同事务时,需要明确协调方案。
数据库事务:写业务结果 + 写 outbox 意图 → 一起提交
发布者:取得未投递意图 → 发送 → 记录已投递
消费者:校验消息 → 在效果边界处理重复 → 保存结果 → 确认这是事务发件箱教学流程。发布后记录前仍可能崩溃,导致重投;它保留发布意图,消费者还需要处理重复。对于同一数据库里的队列,也要确认入队 API 是否真的接入业务事务,名称相同不能证明原子性。
可把发件箱做成业务数据库中的公开应用表。下面是简化 schema:payload 保存可重放输入,时间与投递状态由发布者维护。实际事件 schema、数据保留和索引按流量安排。
CREATE TABLE outbox_events (
id uuid PRIMARY KEY,
topic text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
dispatched_at timestamptz
);
CREATE INDEX outbox_pending ON outbox_events (created_at, id)
WHERE dispatched_at IS NULL;业务写入和 INSERT INTO outbox_events ... 必须用同一个数据库连接和事务。发布者可以在短事务中按顺序取得一条待投递记录,FOR UPDATE SKIP LOCKED 让其他发布者跳过已锁行,调用队列后再记投递时间:
import { Pool } from 'pg';
import { Queue } from 'bullmq';
import { z } from 'zod';
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const queue = new Queue('events', {
connection: { host: '127.0.0.1', port: 6379 },
});
const eventSchema = z.object({
id: z.uuid(), topic: z.string().min(1), payload: z.unknown(),
});
export async function publishOne() {
const db = await pool.connect();
try {
await db.query('BEGIN');
const rows = await db.query(`
SELECT id, topic, payload FROM outbox_events
WHERE dispatched_at IS NULL ORDER BY created_at, id
FOR UPDATE SKIP LOCKED LIMIT 1
`);
if (rows.rows.length === 0) { await db.query('COMMIT'); return false; }
const event = eventSchema.parse(rows.rows[0]);
await queue.add(event.topic, event.payload, { jobId: event.id });
await db.query('UPDATE outbox_events SET dispatched_at = now() WHERE id = $1', [event.id]);
await db.query('COMMIT');
return true;
} catch (error) {
await db.query('ROLLBACK');
throw error;
} finally {
db.release();
}
}这是一条发布的调用链,轮询调度、配置启动校验、连接关闭、告警与退避由宿主接入。示例在等待队列时持有数据库锁,需限制外部等待;高负载可改为有期限的租约领取,并处理租约过期和旧领取者。队列已接受但 COMMIT 失败会留下待投递行,下一次再投递;即使 jobId 减少部分重复,消费者仍按业务身份处理重复。以上未运行新服务。
jobId 提供哪种去重
BullMQ 自定义 jobId 在一个队列内辨认已有任务。任务删除后,同 ID 可以再次入队;不同队列也可使用相同 ID。它不保证业务效果永远只发生一次。BullMQ Job IDs
效果幂等需要持久的业务身份、参数与结果政策。内部数据库写入可以把重复记录和效果放在同一事务里;外部发送、扣款或存储需要对应接收方幂等与查证能力。先记录“已处理”再产生效果会丢工作,先产生效果再记记录则可能重复。
重试保留同一逻辑意图,按可恢复原因退避与抖动,并限制尝试次数及总期限。毒消息进入可观察的失败处理,保留输入身份和原因。自动清理任务历史要与迟到重试、排障和业务记录保留政策相容。
worker 运行在可持续工作的独立进程。下面执行一个可安全重算的摘要任务,展示验证、处理结果、错误和停机接线;它没有数据库或外部副作用,不能据此推导邮件或扣款幂等:
import { createHash } from 'node:crypto';
import { Worker } from 'bullmq';
import { z } from 'zod';
const inputSchema = z.object({ text: z.string().max(100_000) });
const worker = new Worker('events', async job => {
if (job.name !== 'text-digest') throw new Error('Unsupported job topic');
const { text } = inputSchema.parse(job.data);
return { sha256: createHash('sha256').update(text, 'utf8').digest('hex') };
}, {
connection: { host: '127.0.0.1', port: 6379, maxRetriesPerRequest: null },
concurrency: 4,
});
worker.on('error', () => console.error('Worker connection error'));
worker.on('failed', (job) => console.error('Task failed', { id: job?.id }));
process.once('SIGTERM', () => {
void worker.close().catch(() => { process.exitCode = 1; });
});生产者为相应事件写 topic = 'text-digest' 和 { text } payload;重试次数与 backoff 在队列添加选项中明确,失败历史按保留政策处理。worker.close() 等待在途任务,平台仍需给足停机期限;强行终止后租约恢复可能再执行。Redis 地址、认证、TLS、持久化和 noeviction 按部署配置替换,示例本地地址不构成生产配置。BullMQ 生产注意事项
消息怎样接到服务
命令要求执行,事件描述已成立事实,请求响应还需要关联返回和超时。消息信封可以保存事件 ID、schema 版本、关联与原因 ID、发生时刻;每个字段按真实生成和传播规则解释。
NATS core 与 JetStream 的持久化和消费者语义不同;持久 consumer、确认和重投仍要按真实配置核对。RabbitMQ、Kafka 和 gRPC 面向不同运输或请求关系,不能因框架统一 transporter 接口而当成等价实现。NATS JetStream
跨服务超时留下结果不确定性。补偿是新的业务动作,有权限和失败;退款、撤销和回收无法自动恢复原始世界。流程记录各步完成与待处理,使恢复者能够继续或补偿。
长连接怎样交付持续结果
SSE 由服务器向浏览器持续发事件,WebSocket 支持双向消息,Socket.IO 在相应运输上增加事件、重连等能力。Socket.IO 客户端不能直接当成任意 WebSocket 客户端,运输选择也影响代理、平台和回退。
下面是 Next App Router 的公开心跳接口,接到 Web Request/Response 与流。它只展示每条连接的资源所有权,不包含业务通知和持久重放。运行平台必须支持持续响应;静态文件部署无法执行该接口。
export function GET(request: Request) {
const encoder = new TextEncoder();
let timer: ReturnType<typeof setInterval> | undefined;
let ended = false;
let abort = () => {};
const cleanup = () => {
if (timer !== undefined) clearInterval(timer);
request.signal.removeEventListener('abort', abort);
};
const stream = new ReadableStream<Uint8Array>({
start(controller) {
abort = () => {
if (ended) return;
ended = true;
cleanup();
controller.close();
};
if (request.signal.aborted) { abort(); return; }
request.signal.addEventListener('abort', abort, { once: true });
const send = () => {
if (ended) return;
controller.enqueue(encoder.encode(
`event: heartbeat\ndata: ${JSON.stringify({ at: Date.now() })}\n\n`,
));
};
send();
timer = setInterval(send, 15_000);
},
cancel() { ended = true; cleanup(); },
});
return new Response(stream, { headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache, no-transform',
} });
}浏览器以 new EventSource('/api/heartbeat') 取得连接,用 addEventListener('heartbeat', handler) 订阅,在 owner 卸载时 removeEventListener 并 close()。收到 data 后解析 JSON/schema,错误事件更新连接反馈。15 秒只是一项示例间隔,代理缓冲、空闲期限和托管执行时限需按部署设置;片段未运行新 Next 服务。Next Route Handler
业务连接在建立、加入 room 和每次受保护操作时检查会话与对象范围。多实例广播需要对应 adapter 或总线,进程内 emit 只到该进程拥有的连接。数据库提交后直接 emit 仍有丢通知窗口;需要可靠传播时持久保存意图,客户端重连后凭版本或事件进度补齐。SSE 的 Last-Event-ID 只有在服务器保存并实现重放规则时才有恢复作用。
通知可令查询缓存失效,失效后重取主要事实;直接更新缓存需确认消息完整、版本新鲜和权限适用。presence 是短暂观察,持久文档和业务事实有自己的存储规则,不能用“连接在线”代表已经同步。
定时与停机
调度选择时区、错过时刻、重复和重叠政策,多实例调度需共同身份或协调。单进程 protect 或互斥标志不能限制其他进程。业务周期可形成持久身份,避免仅以当前毫秒时间生成每次任务。
停机先停止接收新工作,等待或中止在途,按实际结果确认,再关闭连接。任务租约过期、进程崩溃和重复执行需要相应恢复测试。归档中队列与消息的代码是当时教学/约定片段,本页未运行新队列或服务实验。
最后更新于