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

如何在RabbitMQ服务中用JavaScript async/await执行Shell脚本

在Node.js中结合async/await执行Shell脚本的正确姿势

好问题!你当前的代码里有个关键问题:await并没有真正等待Shell脚本执行完成。原因是Node.js原生的exec是基于回调的API,它不会返回Promise,所以await对它完全不起作用——代码会直接跳过exec的回调逻辑,根本等不到脚本跑完就继续往下走了。

咱们来一步步修正这个问题,同时优化代码的安全性和可读性:

1. 把回调式exec转成Promise风格

Node.js自带的util.promisify工具可以轻松把回调风格的函数转换成支持async/await的Promise版本。首先在文件顶部引入相关模块:

const { promisify } = require('util');
const exec = promisify(require('child_process').exec);

2. 修改RabbitMQ消息处理逻辑

现在可以用await正确等待脚本执行完成,同时用try/catch捕获错误(代替原来回调里的error !== null判断):

console.log(' [*] Waiting for messages in %s. To exit press CTRL+C', q);

channel.consume(q, async function reply(msg) {
  if (!msg) return; // 防御性判断,避免空消息导致报错
  const mongodbUserId = msg.content.toString();
  console.log(' [x] Received %s', mongodbUserId);

  try {
    // 用await等待脚本执行完成,解构获取输出内容
    const { stdout, stderr } = await exec(`./new_user_run_athena.sh ${mongodbUserId}`);
    console.log('Running Athena...');
    console.log('stdout: ' + stdout);
    console.log('stderr: ' + stderr);
  } catch (error) {
    console.log('exec error:', error);
  }

  // 重要:手动确认消息已处理,避免RabbitMQ重复投递未确认的消息
  channel.ack(msg);
});

3. 更安全的选择:规避Shell注入风险

直接把用户输入(比如mongodbUserId)拼接到Shell命令里存在Shell注入风险——如果mongodbUserId包含;、&&这类特殊字符,可能会被执行恶意命令。更安全的做法是用child_process.spawn,它会单独传递参数,不会被解析成Shell命令的一部分:

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

// ...

channel.consume(q, async function reply(msg) {
  if (!msg) return;
  const mongodbUserId = msg.content.toString();
  console.log(' [x] Received %s', mongodbUserId);

  try {
    await new Promise((resolve, reject) => {
      // 用spawn执行脚本,参数单独放在数组中传递
      const childProcess = spawn('./new_user_run_athena.sh', [mongodbUserId]);

      // 实时监听脚本输出
      childProcess.stdout.on('data', (data) => {
        console.log(`stdout: ${data}`);
      });
      childProcess.stderr.on('data', (data) => {
        console.error(`stderr: ${data}`);
      });

      // 监听进程结束事件
      childProcess.on('close', (code) => {
        if (code === 0) {
          resolve();
        } else {
          reject(new Error(`脚本执行失败,退出码:${code}`));
        }
      });
    });
    console.log('Running Athena... 执行完成');
  } catch (error) {
    console.log('exec error:', error);
  }

  channel.ack(msg);
});

额外小提醒

  • 确保./new_user_run_athena.sh拥有可执行权限,可以用chmod +x new_user_run_athena.sh命令设置。
  • 如果你的RabbitMQ消费者需要高可靠性,务必处理消息确认(channel.ack),避免进程崩溃后消息丢失或重复投递。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:02:42