You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Node.js集群开发疑问:Master节点通信及Worker故障转移实现

Node.js集群Master节点通信与Worker故障转移实现方案

一、Master节点间通信实现

方案1:基于共享PostgreSQL数据库的消息传递

利用已有的共享数据库实现跨机器/本地Master间通信,步骤如下:

  1. 创建消息存储表
    在数据库中新增master_messages表,用于存储Master间的消息:
CREATE TABLE master_messages (
    id SERIAL PRIMARY KEY,
    sender_pid INT NOT NULL,
    receiver_pid INT, -- 为NULL时表示广播消息
    content TEXT NOT NULL,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    is_read BOOLEAN DEFAULT FALSE
);

同时新增master表维护Master状态:

CREATE TABLE master (
    master_id INT PRIMARY KEY,
    status VARCHAR(20) DEFAULT 'Online',
    heartbeat TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
  1. Master注册与消息处理逻辑
    在现有Master代码中添加注册、心跳更新、消息轮询与发送逻辑:
if (cluster.isMaster) {
    // 注册当前Master到数据库
    pool.query(
        'INSERT INTO master (master_id) VALUES ($1) ON CONFLICT (master_id) DO UPDATE SET status = $2',
        [process.pid, 'Online'],
        (err) => err && console.error('Master注册失败:', err)
    );

    // 每5秒更新心跳
    setInterval(() => {
        pool.query('UPDATE master SET heartbeat = CURRENT_TIMESTAMP WHERE master_id = $1', [process.pid]);
    }, 5000);

    // 每2秒轮询未读消息
    setInterval(() => {
        pool.query(
            'SELECT * FROM master_messages WHERE (receiver_pid = $1 OR receiver_pid IS NULL) AND is_read = FALSE',
            [process.pid],
            (err, result) => {
                if (err) return console.error('读取消息失败:', err);
                result.rows.forEach(msg => {
                    console.log(`[Master ${process.pid}] 收到Master ${msg.sender_pid} 的消息: ${msg.content}`);
                    // 标记消息为已读
                    pool.query('UPDATE master_messages SET is_read = TRUE WHERE id = $1', [msg.id]);
                });
            }
        );
    }, 2000);

    // 发送消息方法(支持定向/广播)
    function sendMasterMsg(receiverPid, content) {
        pool.query(
            'INSERT INTO master_messages (sender_pid, receiver_pid, content) VALUES ($1, $2, $3)',
            [process.pid, receiverPid, content],
            (err) => err && console.error('发送消息失败:', err)
        );
    }

    // 示例:启动后广播消息
    setTimeout(() => sendMasterMsg(null, `Master ${process.pid} 已启动`), 3000);

    // 原有Master逻辑...
}

方案2:基于UDP广播的本地通信

如果集群部署在同一机器,UDP广播更高效,无需依赖数据库:

const dgram = require('dgram');
const BROADCAST_PORT = 3000;
const BROADCAST_ADDR = '255.255.255.255';

if (cluster.isMaster) {
    const socket = dgram.createSocket('udp4');
    socket.bind(BROADCAST_PORT, () => socket.setBroadcast(true));

    // 接收消息
    socket.on('message', (msg, rinfo) => {
        if (rinfo.port !== BROADCAST_PORT) return; // 过滤自身发送的消息
        console.log(`[Master ${process.pid}] 收到消息: ${msg.toString()}`);
    });

    // 发送广播消息
    function broadcastMsg(content) {
        const msg = Buffer.from(`Master ${process.pid}: ${content}`);
        socket.send(msg, 0, msg.length, BROADCAST_PORT, BROADCAST_ADDR);
    }

    // 每5秒发送心跳广播
    setInterval(() => broadcastMsg('心跳'), 5000);

    // 原有Master逻辑...
}

二、Worker检测Master存活并转为Master

实现逻辑

  1. Worker启动时通过环境变量获取对应Master的PID
  2. 定期查询数据库中Master的心跳时间,超过阈值判定为离线
  3. 检测到Master离线后,以Master身份重新启动自身进程

代码整合修改

在现有Worker逻辑中添加检测与切换逻辑:

else {
    const masterPid = process.env.MASTER_PID;
    const CHECK_INTERVAL = 5000;
    const TIMEOUT_THRESHOLD = 10000; // 10秒超时

    // 定期检测Master状态
    const checkTimer = setInterval(async () => {
        try {
            const result = await pool.query('SELECT heartbeat FROM master WHERE master_id = $1', [masterPid]);
            if (result.rows.length === 0) {
                switchToMaster();
                return;
            }
            const lastHeartbeat = new Date(result.rows[0].heartbeat);
            if (new Date() - lastHeartbeat > TIMEOUT_THRESHOLD) {
                switchToMaster();
            }
        } catch (err) {
            console.error(`[Worker ${process.pid}] 检测Master状态失败:`, err);
        }
    }, CHECK_INTERVAL);

    // 切换为Master进程
    function switchToMaster() {
        console.log(`[Worker ${process.pid}] Master ${masterPid} 离线,切换为Master模式`);
        clearInterval(checkTimer);
        require('child_process').spawn(process.execPath, process.argv, {
            stdio: 'inherit',
            env: { ...process.env, IS_MASTER: 'true' }
        });
        process.exit(0);
    }

    // 原有Worker逻辑
    process.send(`Worker ${process.pid} 已启动`);
}

同时修改Master启动Worker时传递环境变量:

// 替换原有cluster.fork()
cluster.fork({ MASTER_PID: process.pid });

完整修改后代码示例

const { Pool } = require('pg');
const cluster = require('cluster');
const numCPUs = require('os').cpus().length;
const dgram = require('dgram');
const { spawn } = require('child_process');

const dbParams = {
    user: 'postgres',
    host: 'localhost',
    database: 'clustersdb',
    password: 'root',
    port: 5433,
};
const pool = new Pool(dbParams);
const workerPool = new Pool(dbParams);

// UDP广播配置(本地Master通信)
const BROADCAST_PORT = 3000;
const BROADCAST_ADDR = '255.255.255.255';

// 强制启动为Master的标识
const forceMaster = process.env.IS_MASTER === 'true';

if (cluster.isMaster || forceMaster) {
    // 重置cluster身份(针对Worker切换的场景)
    if (forceMaster) {
        cluster.isMaster = true;
        cluster.isWorker = false;
    }

    // Master注册与心跳维护
    pool.query(
        'INSERT INTO master (master_id) VALUES ($1) ON CONFLICT (master_id) DO UPDATE SET status = $2',
        [process.pid, 'Online'],
        (err) => err && console.error('Master注册失败:', err)
    );
    setInterval(() => {
        pool.query('UPDATE master SET heartbeat = CURRENT_TIMESTAMP WHERE master_id = $1', [process.pid]);
    }, 5000);

    // UDP广播通信
    const socket = dgram.createSocket('udp4');
    socket.bind(BROADCAST_PORT, () => socket.setBroadcast(true));
    socket.on('message', (msg, rinfo) => {
        if (rinfo.port !== BROADCAST_PORT) return;
        console.log(`[Master ${process.pid}] 收到消息: ${msg.toString()}`);
    });
    const broadcastMsg = (content) => {
        const msg = Buffer.from(`Master ${process.pid}: ${content}`);
        socket.send(msg, 0, msg.length, BROADCAST_PORT, BROADCAST_ADDR);
    };

    // 启动Worker
    for (let i = 0; i < numCPUs - 2; i++) {
        cluster.fork({ MASTER_PID: process.pid });
    }

    // 原有Master事件监听
    cluster.on('fork', (worker) => {
        insert(worker.process.pid, process.pid);
        broadcastMsg(`启动Worker ${worker.process.pid}`);
    });
    cluster.on('exit', (worker) => {
        update(worker.process.pid);
        const newWorker = cluster.fork({ MASTER_PID: process.pid });
        broadcastMsg(`Worker ${worker.process.pid} 离线,重启为 ${newWorker.process.pid}`);
    });
    cluster.on('message', (worker, message) => {
        console.log(`[Master ${process.pid}] 收到Worker ${worker.process.pid} 消息: ${message}`);
    });
} else {
    const masterPid = process.env.MASTER_PID;
    const CHECK_INTERVAL = 5000;
    const TIMEOUT_THRESHOLD = 10000;

    // Master存活检测
    const checkTimer = setInterval(async () => {
        try {
            const result = await pool.query('SELECT heartbeat FROM master WHERE master_id = $1', [masterPid]);
            if (result.rows.length === 0 || new Date() - new Date(result.rows[0].heartbeat) > TIMEOUT_THRESHOLD) {
                switchToMaster();
            }
        } catch (err) {
            console.error(`[Worker ${process.pid}] 检测Master状态失败:`, err);
        }
    }, CHECK_INTERVAL);

    // 切换为Master
    const switchToMaster = () => {
        console.log(`[Worker ${process.pid}] Master ${masterPid} 离线,切换为Master模式`);
        clearInterval(checkTimer);
        spawn(process.execPath, process.argv, {
            stdio: 'inherit',
            env: { ...process.env, IS_MASTER: 'true' }
        });
        process.exit(0);
    };

    process.send(`Worker ${process.pid} 已启动`);
}

// 原有数据库操作函数
function insert(worker, master) {
    workerPool.query(
        'INSERT INTO worker (worker_id,master_id,status) VALUES ($1, $2, $3)',
        [worker, master, 'Online'],
        (err) => err && console.error('插入Worker数据失败:', err)
    );
}
function update(worker) {
    workerPool.query(
        'UPDATE worker SET status = $1 WHERE worker_id = $2',
        ['Offline', worker],
        (err) => err && console.error('更新Worker状态失败:', err)
    );
}
function fetch(master) {
    workerPool.query('SELECT * FROM worker WHERE master_id = $1', [master], (err, result) => {
        if (err) console.error('查询Worker数据失败:', err);
        else console.log(result.rows.length > 0 ? result.rows : `Master ${master} 无关联Worker`);
    });
}

内容的提问来源于stack exchange,提问作者azerty

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 17:42:04