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

如何通过Node.js调用Python脚本并实现持续交互与结构复用

Node.js 调用长期运行的 Python 守护进程(复用预加载结构)

我来给你梳理一个完整的解决方案,刚好匹配你的需求:让Python进程启动时一次性加载所有复杂结构,之后长期运行等待Node的任务指令;收到任务后用multiprocessing创建子进程执行,完成后返回结果并继续监听,还支持手动终止Python进程。

整体架构思路

  • 启动单长期运行的Python守护进程,进程初始化时加载所有复杂结构(比如AI模型、大型数据集),这一步只做一次
  • Node.js通过**标准输入输出(IPC)**向Python进程发送任务指令(用JSON序列化,跨语言解析更方便)
  • Python主进程收到指令后,用multiprocessing创建子进程执行具体任务——这样主进程不会被任务阻塞,能持续监听新的指令
  • 子进程执行完成后,将结果返回给主进程,主进程再把结果发回Node,之后回到监听状态等待下一个任务

1. Python 守护进程实现

下面是Python脚本的完整代码,我会逐段解释:

import sys
import json
import multiprocessing
from time import sleep

# --------------------------
# 预加载复杂结构(仅进程启动时执行一次)
# --------------------------
print("[Python] 开始加载复杂结构...", file=sys.stderr)
# 这里替换成你的实际复杂结构(比如模型加载、数据集读取)
complex_structure = {
    "model_weights": "模拟的大模型权重数据(实际场景中可能是GB级别的模型)",
    "lookup_table": {str(i): i*2 for i in range(100000)}  # 模拟大型查找表
}
print("[Python] 复杂结构加载完成,等待任务...", file=sys.stderr)

# --------------------------
# 子进程任务执行函数
# --------------------------
def execute_task(task_data):
    """执行具体业务任务,直接复用预加载的complex_structure"""
    task_type = task_data.get("type")
    payload = task_data.get("payload")

    # 模拟任务执行耗时(比如模型推理、数据计算)
    sleep(1)

    # 根据任务类型处理
    if task_type == "calculate":
        num = payload.get("number")
        # 复用预加载的lookup_table做计算
        result = complex_structure["lookup_table"].get(str(num), "未找到对应值")
        return {"status": "success", "result": result}
    elif task_type == "predict":
        input_text = payload.get("text")
        # 复用预加载的model_weights做模拟预测
        result = f"基于预加载模型对'{input_text}'的预测结果:模拟输出"
        return {"status": "success", "result": result}
    else:
        return {"status": "error", "message": "未知的任务类型"}

# --------------------------
# 主进程监听与任务分发逻辑
# --------------------------
def main():
    while True:
        # 读取Node发来的一行JSON消息(必须换行分隔)
        line = sys.stdin.readline()
        if not line:
            # 输入流关闭(比如Node终止进程),退出循环
            print("[Python] 输入流关闭,准备退出守护进程", file=sys.stderr)
            break

        try:
            task = json.loads(line.strip())
            print(f"[Python] 收到新任务: {task}", file=sys.stderr)

            # 用进程池创建子进程执行任务(避免阻塞主进程监听)
            with multiprocessing.Pool(1) as pool:
                task_result = pool.apply(execute_task, (task,))

            # 将结果转为JSON,通过stdout返回给Node(必须flush确保立即发送)
            print(json.dumps(task_result), flush=True)
            print(f"[Python] 任务完成,已返回结果: {task_result}", file=sys.stderr)

        except json.JSONDecodeError:
            error_msg = {"status": "error", "message": "收到无效的JSON格式消息"}
            print(json.dumps(error_msg), flush=True)
            print("[Python] 错误:无法解析收到的消息,不是有效JSON", file=sys.stderr)
        except Exception as e:
            error_msg = {"status": "error", "message": str(e)}
            print(json.dumps(error_msg), flush=True)
            print(f"[Python] 任务执行出错: {str(e)}", file=sys.stderr)

if __name__ == "__main__":
    main()

代码关键点说明

  • 预加载逻辑:放在进程启动最开始,只执行一次,后续所有任务都复用这个结构
  • 日志隔离:用sys.stderr打印日志,避免和stdout的结果数据混淆(Node会把Python的stderr输出到自己的控制台)
  • 子进程管理:用multiprocessing.Pool创建子进程,自动处理进程的创建和销毁;如果任务量很大,可以设置更大的进程池大小(比如Pool(4))
  • 消息格式:用JSON序列化任务和结果,确保跨语言数据结构一致;每次发送/读取都用换行分隔,避免粘包问题

2. Node.js 客户端实现

下面是Node.js的代码,负责启动Python守护进程、发送任务、接收结果和手动终止进程:

const { spawn } = require('child_process');
const path = require('path');

// 启动Python守护进程(替换成你的Python脚本路径)
const pythonDaemon = spawn('python', [path.join(__dirname, 'python_daemon.py')], {
    stdio: ['pipe', 'pipe', 'inherit'] // stdin/stdout用于IPC通信,stderr继承到Node控制台
});

// 监听Python进程的stdout,接收任务结果
pythonDaemon.stdout.on('data', (data) => {
    try {
        const result = JSON.parse(data.toString().trim());
        console.log('[Node] 收到Python返回结果:', result);
    } catch (err) {
        console.error('[Node] 解析Python返回结果出错:', err);
    }
});

// 监听Python进程的退出事件
pythonDaemon.on('exit', (code) => {
    console.log(`[Node] Python守护进程已退出,退出码: ${code}`);
});

// 向Python发送任务的工具函数
function sendTask(task) {
    // 必须加换行符,因为Python用readline读取一行作为一个任务
    const taskJson = JSON.stringify(task) + '\n';
    pythonDaemon.stdin.write(taskJson, (err) => {
        if (err) {
            console.error('[Node] 发送任务失败:', err);
        } else {
            console.log('[Node] 任务已发送:', task);
        }
    });
}

// --------------------------
// 示例:发送测试任务
// --------------------------
// 1秒后发送计算任务
setTimeout(() => {
    sendTask({
        type: 'calculate',
        payload: { number: 123 }
    });
}, 1000);

// 3秒后发送预测任务
setTimeout(() => {
    sendTask({
        type: 'predict',
        payload: { text: 'Hello Python Daemon' }
    });
}, 3000);

// --------------------------
// 手动终止Python进程的函数
// --------------------------
function terminateDaemon() {
    // 先关闭stdin,让Python进程退出监听循环
    pythonDaemon.stdin.end();
    // 强制杀死进程(确保彻底退出)
    pythonDaemon.kill();
    console.log('[Node] 已手动终止Python守护进程');
}

// 示例:5秒后自动终止进程
setTimeout(terminateDaemon, 5000);

代码关键点说明

  • 进程启动:用spawn而不是exec,因为spawn适合长期运行的进程,能实时处理IO流
  • IPC配置:stdio: ['pipe', 'pipe', 'inherit']让Node和Python的stdin/stdout互通,stderr直接输出到Node控制台,方便调试
  • 任务发送:必须给JSON字符串加换行符,否则Python的readline会一直等待输入
  • 终止逻辑:先关闭stdin让Python优雅退出,再用kill确保进程彻底终止

3. 关键注意事项

  • 子进程内存优化:如果你的复杂结构非常大(比如几GB的模型),子进程复制父进程内存会有开销。这种情况下可以考虑用multiprocessing.Manager共享只读数据,或者把结构放在外部存储(比如Redis)中,子进程按需读取
  • 错误处理增强:可以给任务添加唯一ID,方便追踪任务的请求和响应;同时处理Python进程崩溃的情况,比如在Node中监听error事件,自动重启Python进程
  • 性能优化:如果任务并发量高,Python端可以用固定大小的进程池,避免频繁创建销毁子进程;Node端可以用任务队列控制发送频率,避免Python进程过载
  • 安全问题:如果任务来自外部输入,要做参数校验,避免注入攻击;同时限制Python进程的资源使用(比如CPU、内存),防止影响系统稳定性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:11:12