分布式系统下的JavaScript消息队列实现需要兼顾消息可靠性、服务解耦、跨节点通信等核心需求,通常可以基于Node.js的网络通信能力和持久化存储来构建基础实现框架。

核心设计要点
分布式场景下的消息队列和单机版本有明显区别,需要重点考虑以下几个维度:
- 消息持久化:避免节点宕机导致消息丢失,需要将消息存储到可靠的存储介质中
- 集群同步:多个队列节点之间需要同步消息状态,保证消费者可以跨节点获取消息
- 消费确认:消费者处理完消息后需要返回确认信号,避免消息重复消费或丢失
- 负载均衡:多个消费者节点之间合理分配消息,提升整体处理效率
基础实现示例
以下示例基于Node.js的net模块实现节点间通信,使用SQLite做本地消息持久化,实现简单的分布式消息队列核心功能。
消息队列服务端实现
服务端负责接收生产者发送的消息、存储消息、向消费者推送消息,同时处理节点间的状态同步。
const net = require('net');
const sqlite3 = require('sqlite3').verbose();
const crypto = require('crypto');
// 初始化本地消息存储数据库
const db = new sqlite3.Database('./queue.db', (err) => {
if (err) {
console.error('数据库初始化失败', err);
} else {
// 创建消息表,存储消息内容、状态、所属队列等信息
db.run(`CREATE TABLE IF NOT EXISTS messages (
id TEXT PRIMARY KEY,
queue_name TEXT NOT NULL,
content TEXT NOT NULL,
status TEXT DEFAULT 'pending',
create_time INTEGER NOT NULL,
consumer_id TEXT
)`);
// 创建节点信息表,存储集群内的其他节点地址
db.run(`CREATE TABLE IF NOT EXISTS nodes (
node_id TEXT PRIMARY KEY,
host TEXT NOT NULL,
port INTEGER NOT NULL
)`);
}
});
// 生成唯一消息ID
function generateMsgId() {
return crypto.randomBytes(16).toString('hex');
}
// 消息队列服务端类
class QueueServer {
constructor(port, host = '0.0.0.0') {
this.port = port;
this.host = host;
this.nodeId = generateMsgId();
this.consumers = new Map(); // 存储注册的消费者信息
this.server = net.createServer(this.handleConnection.bind(this));
}
// 启动服务端
start() {
this.server.listen(this.port, this.host, () => {
console.log(`队列节点 ${this.nodeId} 启动成功,监听 ${this.host}:${this.port}`);
});
}
// 处理客户端连接
handleConnection(socket) {
socket.on('data', (data) => {
try {
const request = JSON.parse(data.toString());
this.handleRequest(socket, request);
} catch (e) {
socket.write(JSON.stringify({ code: 400, msg: '请求格式错误' }));
}
});
socket.on('error', (err) => {
console.error('客户端连接错误', err);
});
}
// 处理不同类型的请求
handleRequest(socket, request) {
switch (request.type) {
case 'produce': // 生产消息
this.handleProduce(socket, request);
break;
case 'register_consumer': // 注册消费者
this.handleRegisterConsumer(socket, request);
break;
case 'pull_message': // 拉取消息
this.handlePullMessage(socket, request);
break;
case 'ack_message': // 确认消息消费
this.handleAckMessage(socket, request);
break;
case 'register_node': // 注册集群节点
this.handleRegisterNode(socket, request);
break;
default:
socket.write(JSON.stringify({ code: 404, msg: '未知请求类型' }));
}
}
// 处理生产消息请求
handleProduce(socket, request) {
const { queueName, content } = request;
if (!queueName || !content) {
socket.write(JSON.stringify({ code: 400, msg: '队列名称和消息内容不能为空' }));
return;
}
const msgId = generateMsgId();
const createTime = Date.now();
// 持久化消息到本地数据库
db.run(
'INSERT INTO messages (id, queue_name, content, create_time) VALUES (?, ?, ?, ?)',
[msgId, queueName, content, createTime],
(err) => {
if (err) {
socket.write(JSON.stringify({ code: 500, msg: '消息存储失败' }));
} else {
// 通知对应队列的消费者有新消息
this.notifyConsumers(queueName);
socket.write(JSON.stringify({ code: 200, msg: '消息发送成功', msgId }));
}
}
);
}
// 注册消费者
handleRegisterConsumer(socket, request) {
const { queueName, consumerId } = request;
if (!queueName || !consumerId) {
socket.write(JSON.stringify({ code: 400, msg: '队列名称和消费者ID不能为空' }));
return;
}
if (!this.consumers.has(queueName)) {
this.consumers.set(queueName, new Set());
}
this.consumers.get(queueName).add({ consumerId, socket });
socket.write(JSON.stringify({ code: 200, msg: '消费者注册成功' }));
}
// 拉取消息
handlePullMessage(socket, request) {
const { queueName, consumerId } = request;
// 查询待处理的消息
db.get(
'SELECT * FROM messages WHERE queue_name = ? AND status = "pending" ORDER BY create_time ASC LIMIT 1',
[queueName],
(err, row) => {
if (err) {
socket.write(JSON.stringify({ code: 500, msg: '查询消息失败' }));
} else if (row) {
// 标记消息为处理中,绑定消费者ID
db.run(
'UPDATE messages SET status = "processing", consumer_id = ? WHERE id = ?',
[consumerId, row.id],
(updateErr) => {
if (updateErr) {
socket.write(JSON.stringify({ code: 500, msg: '消息状态更新失败' }));
} else {
socket.write(JSON.stringify({ code: 200, msg: '拉取消息成功', data: row }));
}
}
);
} else {
socket.write(JSON.stringify({ code: 200, msg: '暂无待处理消息', data: null }));
}
}
);
}
// 确认消息消费
handleAckMessage(socket, request) {
const { msgId, consumerId } = request;
// 更新消息状态为已完成
db.run(
'UPDATE messages SET status = "finished" WHERE id = ? AND consumer_id = ?',
[msgId, consumerId],
(err) => {
if (err) {
socket.write(JSON.stringify({ code: 500, msg: '消息确认失败' }));
} else {
socket.write(JSON.stringify({ code: 200, msg: '消息确认成功' }));
}
}
);
}
// 注册集群节点
handleRegisterNode(socket, request) {
const { nodeId, host, port } = request;
db.run(
'INSERT OR REPLACE INTO nodes (node_id, host, port) VALUES (?, ?, ?)',
[nodeId, host, port],
(err) => {
if (err) {
socket.write(JSON.stringify({ code: 500, msg: '节点注册失败' }));
} else {
socket.write(JSON.stringify({ code: 200, msg: '节点注册成功' }));
}
}
);
}
// 通知对应队列的消费者有新消息
notifyConsumers(queueName) {
const consumers = this.consumers.get(queueName);
if (consumers) {
consumers.forEach(consumer => {
try {
consumer.socket.write(JSON.stringify({ type: 'new_message', queueName }));
} catch (e) {
// 移除失效的消费者连接
consumers.delete(consumer);
}
});
}
}
}
// 启动队列节点,端口可以自定义
const queueServer = new QueueServer(3000);
queueServer.start();
生产者客户端实现
生产者负责向队列节点发送消息,实现逻辑相对简单,只需要和队列节点建立TCP连接发送生产请求即可。
const net = require('net');
class MessageProducer {
constructor(nodeHost, nodePort) {
this.nodeHost = nodeHost;
this.nodePort = nodePort;
this.socket = null;
}
// 连接队列节点
connect() {
return new Promise((resolve, reject) => {
this.socket = net.createConnection({ host: this.nodeHost, port: this.nodePort }, () => {
resolve();
});
this.socket.on('error', reject);
});
}
// 发送消息
sendMessage(queueName, content) {
return new Promise((resolve, reject) => {
const request = {
type: 'produce',
queueName,
content
};
this.socket.write(JSON.stringify(request));
this.socket.once('data', (data) => {
try {
const response = JSON.parse(data.toString());
resolve(response);
} catch (e) {
reject(new Error('响应解析失败'));
}
});
});
}
// 关闭连接
close() {
if (this.socket) {
this.socket.end();
}
}
}
// 使用示例
async function testProducer() {
const producer = new MessageProducer('127.0.0.1', 3000);
await producer.connect();
const result = await producer.sendMessage('order_queue', '新订单创建,订单号:20240501001');
console.log('生产消息结果', result);
producer.close();
}
testProducer();
消费者客户端实现
消费者需要先注册到队列节点,然后拉取消息处理,处理完成后发送确认请求。
const net = require('net');
class MessageConsumer {
constructor(nodeHost, nodePort, queueName, consumerId) {
this.nodeHost = nodeHost;
this.nodePort = nodePort;
this.queueName = queueName;
this.consumerId = consumerId;
this.socket = null;
}
// 连接队列节点并注册消费者
connect() {
return new Promise((resolve, reject) => {
this.socket = net.createConnection({ host: this.nodeHost, port: this.nodePort }, () => {
// 发送注册请求
const registerReq = {
type: 'register_consumer',
queueName: this.queueName,
consumerId: this.consumerId
};
this.socket.write(JSON.stringify(registerReq));
this.socket.once('data', (data) => {
try {
const response = JSON.parse(data.toString());
if (response.code === 200) {
resolve();
} else {
reject(new Error(response.msg));
}
} catch (e) {
reject(new Error('注册响应解析失败'));
}
});
});
this.socket.on('error', reject);
});
}
// 开始消费消息
startConsume(messageHandler) {
// 监听服务端的新消息通知
this.socket.on('data', async (data) => {
try {
const notification = JSON.parse(data.toString());
if (notification.type === 'new_message') {
// 拉取消息
await this.pullAndHandleMessage(messageHandler);
}
} catch (e) {
// 忽略非通知类的响应数据
}
});
// 初始拉取一次消息
this.pullAndHandleMessage(messageHandler);
}
// 拉取并处理消息
async pullAndHandleMessage(messageHandler) {
const pullReq = {
type: 'pull_message',
queueName: this.queueName,
consumerId: this.consumerId
};
this.socket.write(JSON.stringify(pullReq));
const data = await new Promise(resolve => {
this.socket.once('data', resolve);
});
try {
const response = JSON.parse(data.toString());
if (response.code === 200 && response.data) {
// 处理消息
await messageHandler(response.data);
// 发送确认请求
const ackReq = {
type: 'ack_message',
msgId: response.data.id,
consumerId: this.consumerId
};
this.socket.write(JSON.stringify(ackReq));
}
} catch (e) {
console.error('消息处理失败', e);
}
}
// 关闭连接
close() {
if (this.socket) {
this.socket.end();
}
}
}
// 使用示例
async function testConsumer() {
const consumer = new MessageConsumer('127.0.0.1', 3000, 'order_queue', 'consumer_001');
await consumer.connect();
consumer.startConsume(async (message) => {
console.log('收到消息', message.content);
// 模拟消息处理耗时
await new Promise(resolve => setTimeout(resolve, 1000));
});
}
testConsumer();
实现优化方向
上述示例是基础实现,实际分布式场景下还需要做更多优化:
- 增加消息重试机制,处理中状态的消息超过一定时间没有确认,重新标记为待处理
- 实现消息分区,不同队列的消息可以分配到不同的节点处理,提升集群吞吐量
- 增加消息TTL机制,过期未消费的消息自动丢弃或转移到死信队列
- 完善集群节点健康检查,自动移除失效节点,保证消息可以路由到可用节点
注意:实际生产环境中建议优先使用成熟的消息队列中间件如RabbitMQ、Kafka等,上述实现仅用于理解分布式消息队列的核心原理,不建议直接用于生产环境。
JavaScript消息队列分布式系统Node.js修改时间:2026-07-23 22:27:59