分布式系统下的JavaScript消息队列如何实现

来源:建站技术作者:广州程序员头衔:程序员
导读:本期聚焦于小伙伴创作的《分布式系统下的JavaScript消息队列如何实现》,敬请观看详情,探索知识的价值。以下视频、文章将为您系统阐述其核心内容与价值。如果您觉得《分布式系统下的JavaScript消息队列如何实现》有用,将其分享出去将是对创作者最好的鼓励。

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

分布式系统下的JavaScript消息队列如何实现

核心设计要点

分布式场景下的消息队列和单机版本有明显区别,需要重点考虑以下几个维度:

  • 消息持久化:避免节点宕机导致消息丢失,需要将消息存储到可靠的存储介质中
  • 集群同步:多个队列节点之间需要同步消息状态,保证消费者可以跨节点获取消息
  • 消费确认:消费者处理完消息后需要返回确认信号,避免消息重复消费或丢失
  • 负载均衡:多个消费者节点之间合理分配消息,提升整体处理效率

基础实现示例

以下示例基于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

免责声明:​ 已尽一切努力确保本网站所含信息的准确性。网站内容多为原创整理与精心编撰,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们处理。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。