消息推送系统几乎是所有中大型应用的标配,不论是站内通知、实时告警还是运营活动的批量触达,核心目标都是将消息可靠、高效地传递给目标用户,同时避免对服务器造成过大压力。在 Node.js 技术栈下实现这样一套系统,需要解决好三个关键环节:批量推送的架构与效率、消息状态的可靠追踪、推送速率的合理控制。
本节不涉及具体的推送通道(如 APNs、FCM、邮件等),而是聚焦于 Node.js 服务端如何设计消息分发、状态管理和流控机制的核心骨架。
26.5.1 批量推送:从逐条发送到流水线处理
当需要向数万、数十万用户推送同一条消息时,一条一条地创建记录并调用推送接口,不仅耗时巨大,还会瞬间拉高数据库和外部服务的压力。批量推送的核心思路是化零为整,并尽可能利用 Node.js 的异步并发能力。
1. 数据库写入的批量化
假设我们的消息队列依托于 MySQL,创建消息记录通常使用单条 INSERT。但在批量推送时,重复发送上万条相似的 INSERT 语句,网络往返和 SQL 解析开销会很高。
合理的做法是使用批量插入:
// 不推荐:循环单条插入
for (const userId of userIds) {
await db.query('INSERT INTO messages (user_id, content, status) VALUES (?, ?, ?)', [userId, content, 'pending']);
}
// 推荐:使用批量插入语法
const values = userIds.map(() => '(?, ?, ?)').join(', ');
const params = userIds.flatMap(uid => [uid, content, 'pending']);
await db.query(`INSERT INTO messages (user_id, content, status) VALUES ${values}`, params);
大多数 ORM 也支持批量创建,例如 Sequelize 的 bulkCreate,Prisma 的 createMany。在具体实现时,还需要考虑单次 SQL 语句的长度和参数数量限制,可以将总体分批成每批 500~1000 条插入。
2. 推送动作的异步并行与任务队列
即使消息记录已入库,实际将通知推送到用户设备(如 WebSocket、App Push)仍然可能比较慢。直接在主请求中批量执行推送会导致接口响应超时,因此更常见的方案是将推送任务异步化。
一个典型设计是:API 收到批量推送请求后,迅速将任务放入消息队列(如 Bull),然后立即返回“推送任务已提交”。后台的 Worker 进程从队列取出任务,逐步处理每一批用户。
// 使用 Bull 队列
const pushQueue = new Bull('push notifications', 'redis://...');
// 接收请求,批量创建消息记录并投放推送任务
app.post('/api/batch-push', async (req, res) => {
const { userIds, content } = req.body;
// 1. 批量写入消息
await Message.bulkCreate(userIds.map(uid => ({ userId: uid, content, status: 'pending' })));
// 2. 将推送任务分组放入队列,避免单任务过大
const batchSize = 1000;
for (let i = 0; i < userIds.length; i += batchSize) {
const batch = userIds.slice(i, i + batchSize);
await pushQueue.add('batch-push', { userIds: batch, content });
}
res.json({ message: '推送任务已提交', total: userIds.length });
});
Worker 进程可以开启多个(Bull 支持并发处理),利用 Node.js 的异步并发能力并行地向用户推送。针对在线用户的 WebSocket 推送,可以在内存中维护用户连接映射,以实现实时下发;离线用户的 App Push 则调用第三方 SDK。
// Worker 处理推送任务
pushQueue.process('batch-push', async (job) => {
const { userIds, content } = job.data;
for (const userId of userIds) {
try {
const ws = connectedClients.get(userId);
if (ws && ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify({ type: 'notification', content }));
await Message.update({ status: 'delivered' }, { where: { userId, content } });
} else {
// 离线推送交给第三方
await pushToFCM(userId, content);
await Message.update({ status: 'sent_to_fcm' }, { where: { userId, content } });
}
} catch (err) {
await Message.update({ status: 'failed', error: err.message }, { where: { userId, content } });
}
}
});
3. 大规模推送的架构优化
当用户量达到百万级,单一 Redis 队列和 Worker 可能成为瓶颈,此时可以:
- 按用户分片:将用户 ID 哈希到多个队列,由不同的 Worker 集群处理。
- 使用流式读取:不再一次性加载所有用户 ID 到内存,而是通过数据库游标或流式查询,边取边处理。
- 聚合推送接口:一些第三方推送服务支持一次 API 调用推送给多个用户(如 FCM 的 multicast),充分利用批量接口减少网络开销。
26.5.2 状态回执:消息流转的可靠追踪
在生产环境中,仅将消息“发出”是不够的,我们需要清楚地知道每一条消息到达了哪个环节,以便于故障排查、数据分析和补偿处理。状态回执系统的设计要点在于状态的精确建模和状态流转的幂等性。
1. 状态定义与数据库设计
消息状态通常会经历以下生命周期:
待发送(pending)→ 已投递(delivered)/ 发送失败(failed)→ 用户已读(read)
表中至少应包含以下字段:
CREATE TABLE push_records (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
message_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
status ENUM('pending', 'delivered', 'read', 'failed') DEFAULT 'pending',
error_message TEXT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_user_message (user_id, message_id)
);
如果一条消息需要向多个用户发送,可以采用 message 主表 + push_records 明细表 的模型。
2. 状态更新的幂等性保障
由于网络波动或重试机制,同一个消息的状态回执可能被多次通知。例如,用户 WebSocket 断开时恰好收到消息,客户端和服务器可能同时尝试更新状态为“delivered”。为了避免数据错乱,状态更新必须设计为单向流转且具备幂等性。
常用方法:
- 在更新时增加前置状态条件:
UPDATE push_records SET status = 'delivered' WHERE id = ? AND status = 'pending'
- 使用乐观锁(版本号) 或 数据库事务 确保并发更新安全。
- 如果使用 Redis 存储临时状态,可以结合 Lua 脚本实现原子操作。
3. 实时回执与对账机制
对于 WebSocket 等在线推送,客户端在收到消息后应立即发送 ACK:
// 服务端监听 ACK
ws.on('message', async (data) => {
const msg = JSON.parse(data);
if (msg.type === 'ack' && msg.messageId) {
await updateStatus(msg.messageId, 'delivered');
// 可触发后续业务流程,如发送已读通知
}
});
然而,客户端可能因为 Bug 或网络问题未发送 ACK,导致消息状态永远停留在 “pending”。因此需要离线对账机制:
- 定时任务每隔一段时间(如 5 分钟)扫描
status='pending'且创建时间超过阈值的记录,检查是否已实际投递(如调用第三方回执查询接口)。 - 对于确认为失败的消息,标记为
failed并触发告警或自动重试。
对账任务的 Node.js 实现可以利用 node-cron 或 Bull 的重复任务:
const cron = require('node-cron');
cron.schedule('*/5 * * * *', async () => {
const staleRecords = await db.query(
"SELECT id FROM push_records WHERE status = 'pending' AND created_at < NOW() - INTERVAL 10 MINUTE LIMIT 1000"
);
for (const record of staleRecords) {
// 查询第三方或重新尝试连接检查
// ...
}
});
26.5.3 限流控制:保护系统的最后防线
消息推送如果瞬间全量下发,不仅可能压垮自己的数据库和推送队列,也可能触发第三方 API 的频率限制,甚至被目标运营商判定为 DDoS。因此,在系统的多个层面实施限流至关重要。
1. 接入层全局限流
首先,管理后台和对外开放的 API 接口需要做请求频率限制,防止人为误操作或恶意调用批量推送接口。
使用 express-rate-limit 或 @nestjs/throttler 可以快速实现:
const rateLimit = require('express-rate-limit');
const pushLimiter = rateLimit({
windowMs: 1 * 60 * 1000, // 1分钟
max: 5, // 最多5个批量推送请求
message: '推送请求太频繁,请稍后再试'
});
app.use('/api/batch-push', pushLimiter);
2. 消息队列层面的消费速率控制
即使接入层有限流,队列中可能仍然堆积了大量待处理任务。Bull 队列支持对 Worker 的并发数和速度进行精细控制,防止下游过载。
const pushQueue = new Bull('push notifications', { redis: { host: '...' } });
// 全局并发限制:最多同时处理 10 个推送任务
pushQueue.process(10, 'batch-push', async (job) => { ... });
// 为队列添加速率限制:每10秒最多处理50个任务
const queueSettings = await pushQueue.getJobCounts();
pushQueue.pause(); // 暂定以便配置
// 也可以在添加任务时设置 delay 控制
更复杂的场景可以结合 bottleneck 库,对特定的第三方 API 调用进行精细限流:
const Bottleneck = require('bottleneck');
const fcmLimiter = new Bottleneck({
maxConcurrent: 10, // 最多同时 10 个 FCM 请求
minTime: 100 // 每个请求间隔最少 100ms
});
async function pushToFCM(userId, content) {
return fcmLimiter.schedule(() => {
return fcmSDK.send({ token: getUserToken(userId), notification: { body: content } });
});
}
3. 用户级别的防打扰控制
从产品角度,频繁给同一个用户推送消息会导致反感甚至卸载。在服务端应当维护用户级别的推送频次控制:
- 利用 Redis 记录每个用户最后一次推送时间戳和每日推送次数。
- 在处理推送任务时进行检查:
async function canPush(userId) {
const todayKey = `push_limit:${userId}:${new Date().toISOString().slice(0,10)}`;
const count = await redis.incr(todayKey);
await redis.expire(todayKey, 86400); // 当天有效
const lastPushKey = `last_push:${userId}`;
const lastPush = await redis.get(lastPushKey);
if (lastPush && Date.now() - Number(lastPush) < 60000) { // 1分钟内不重复推
return false;
}
if (count > 5) { // 每天最多5条
return false;
}
await redis.set(lastPushKey, Date.now());
return true;
}
用户级限流与推送失败不同,它属于正常的策略跳过,我们可以将这类消息状态更新为 skipped,并记录跳过原因,便于后续分析。
4. 全局限流的动态调整与熔断
在高负载或出现第三方服务故障时,系统应该能够自动降级。可以通过监控关键指标(如队列长度、延迟、失败率)来触发限流参数动态调整甚至熔断:
- 当 Bull 队列中
waiting任务数超过阈值时,自动降低 API 接收速率或暂停新任务。 - 使用断路器模式(例如
opossum库)保护第三方推送服务,失败率达到 50% 时熔断一段时间,避免雪崩。
const circuitBreaker = new Opossum(pushToFCM, {
timeout: 10000, // 超时 10s
errorThresholdPercentage: 50,
resetTimeout: 30000
});
circuitBreaker.fallback(() => {
// 熔断时存入死信队列,等待恢复
deadLetterQueue.add({ userId, content });
});
26.5.4 系统整体协作:一张蓝图
总结以上三点,一个完整的消息推送系统通常由以下几个模块协作:
- 管理端 API:接收批量推送请求,快速入库并投放队列任务,受全局限流保护。
- 消息队列(Bull/BullMQ):解耦请求与推送,支持重试、延迟、速率限制。
- 推送 Worker:消费队列任务,区分在线(WebSocket)和离线(三方 SDK)通道执行推送,并更新状态记录。
- 状态数据库:持久化每一条消息的流转状态,提供对账依据。
- 用户限速服务:基于 Redis 实现用户级别推送频控,保证体验。
- 监控与对账:定时任务修复缺失状态,监控队列深度和失败率,动态触发限流调整。
所有组件都基于 Node.js 的异步模型构建,使得整体系统能够以较少的资源支撑大规模的消息分发。在实际开发中,还要根据具体的业务体量选择合适的中件间(Redis、RabbitMQ),并做好充分的压力测试,确保在百万线上用户峰值时系统依然平稳运行。
消息推送看似简单,实则涉及到服务端设计的方方面面,从数据库的批量写入到分布式下的状态一致性,再到流量控制体系的构建。理清这三个核心环节,就能搭建起一套健壮、可扩展的推送骨架,为后期的业务增长铺平道路。