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

使用node-schedule在Cron任务内调度子任务失败求助

问题排查与解决:node-schedule异步查询导致任务中断

问题描述

使用node-schedule实现每5秒运行的Cron任务,逻辑为检查子任务是否已调度,未调度则创建新子任务。但执行await InteractionFilter.findOne()等PostgreSQL异步查询时,后续代码直接终止,子任务无法创建;注释这些异步查询后,子任务能正常调度并被循环识别。

故障版本代码

scheduleJob("worloads cron", "*/5 * * * * *", async function () {
  console.log("This job runs every 5 seconds");

  if (!connection) { // 创建数据库连接(如果不存在)
    Logger.info("Creating workloads DB connection.");
    connection = await createConnection({
      type: "postgres",
      url: process.env.DATABASE_URL,
      logging: true, // todo: 生产环境关闭
    });
    return Promise.resolve();
  } else if (connection && connection.isConnected) {

    // 查询需要调度的数据
    const data: Workload[] = await connection.query(`
       SELECT workload."status"
       from "workload" workload;';
  `);

    console.log("jobs", scheduledJobs);

    for (const workload of data) {
      if (scheduledJobs[String(workload.id)]) {
        console.log(`Job with id ${workload.id} exists!`);
      } else {
        Logger.info(
          `Job with id ${workload.id} does not exist. Creating new scheduled job!`
        );
        // 从此处开始代码不再执行,任务无法创建
        const filter = await InteractionFilter.findOne({
          where: { id: workload.interactionFilter },
        });

        Logger.info("FILTER", filter); // 无日志输出
        const qb = await connection
          .getRepository(Interaction)
          .createQueryBuilder("interaction")
          .leftJoinAndSelect(
            `interaction.${filter.integrationSystem.toLowerCase()}_metadata`,
            "metadata"
          )
          .where("interaction.system = :system", {
            system: filter.integrationSystem,
          })
          .orderBy("interaction.createdAt", "ASC")
          .getMany();

        const filter_values = await connection
          .getRepository(InteractionFilter_Metadata)
          .createQueryBuilder("filtervalue")
          .leftJoinAndSelect("filtervalue.metadata", "metadata")
          .where("filtervalue.filterId = :filterId", {
            filterId: workload.interactionFilter,
          })
          .distinct()
          .getMany();

        Logger.info("FILTER VALUES", filter_values);

        scheduleJob(
          String(workload.id),
          dayjs(workload.executeTime).toDate(),
          function () {
            // todo: 此处创建交互逻辑
            console.log("TEST");
          }
        );
        console.log("SCHEDULE TASK AFTER", scheduledJobs);
      }
    }

    return Promise.resolve();
  }
});


expose(() => true);

故障运行输出

INFO [12-05-2023 16:22:40]: Creating workloads DB connection.
This job runs every 3 seconds
query: 
       SELECT workload."status"
       from "workload" workload;

jobs {
  'worloads crong': Job {

  }
}
INFO [12-05-2023 16:22:50]: Job with id 53 does not exist. Creating new scheduled job!

// 此处本该创建任务,但后续输出显示任务未被识别:

This job runs every 3 seconds
query: 
       SELECT workload."status"
       from "workload" workload;

jobs {
  'worloads crong': Job {
    ...
  }
}
INFO [12-05-2023 16:23:00]: Job with id 53 does not exist. Creating new scheduled job!

可运行版本(注释异步查询后)

scheduleJob("worloads crong", "*/10 * * * * *", async function () {
  console.log("This job runs every 3 seconds");
  if (!connection) {
    Logger.info("Creating workloads DB connection.");
    connection = await createConnection({
      type: "postgres",
      url: process.env.DATABASE_URL,
      logging: true, 
    });
    return Promise.resolve();
  } else if (connection && connection.isConnected) {
            const data: Workload[] = await connection.query(`
       SELECT workload."status"
       from "workload" workload;';
  `);


    console.log("jobs", scheduledJobs);

    for (const workload of data) {
      if (scheduledJobs[String(workload.id)]) {
        console.log(`Job with id ${workload.id} exists!`);
      } else {
        Logger.info(
          `Job with id ${workload.id} does not exist. Creating new scheduled job!`
        );

        scheduleJob(
          String(workload.id),
          dayjs(workload.executeTime).toDate(),
          function () {
            // todo: 此处创建交互逻辑
            console.log("TEST");
          }
        );
        console.log("SCHEDULE TASK AFTER", scheduledJobs);
      }
    }

    return Promise.resolve();
  }
});

可运行版本输出

INFO [12-05-2023 16:27:40]: Creating workloads DB connection.
    This job runs every 3 seconds
    query: 
           SELECT workload."status"
           from "workload" workload;
    
    jobs {
      'worloads crong': Job {
      }
    }
    SCHEDULE TASK AFTER {
      '53': Job {
...
      },
      'worloads crong': Job {
...      
}
    }
    INFO [12-05-2023 16:27:50]: Job with id 53 does not exist. Creating new scheduled job!
    This job runs every 3 seconds
        query: 
           SELECT workload."status"
           from "workload" workload;
    jobs {
      '53': Job {
       ...
      },
      'worloads crong': Job {
       ...
      }
    }
    Job with id 53 exists!

排查结论与解决方案

核心问题

异步查询抛出未捕获的异常,导致整个async任务终止,后续代码无法执行。node-schedule的异步任务如果出现未处理的异常,会直接中断任务流程,不会继续执行剩余逻辑。

解决步骤

  1. 添加异常捕获:给所有异步操作包裹try/catch块,捕获并记录异常,避免任务中断。
  2. 验证数据完整性:确保查询返回的workload包含id、interactionFilter、executeTime等必要字段,避免后续逻辑因数据缺失报错。
  3. 检查实体映射与查询有效性:确认InteractionFilter等实体与数据库表正确映射,且workload.interactionFilter为有效存在的ID,避免查询返回空值或抛出不存在错误。

修改后的代码示例

scheduleJob("worloads cron", "*/5 * * * * *", async function () {
  console.log("This job runs every 5 seconds");

  if (!connection) {
    Logger.info("Creating workloads DB connection.");
    try {
      connection = await createConnection({
        type: "postgres",
        url: process.env.DATABASE_URL,
        logging: true,
      });
    } catch (err) {
      Logger.error("创建数据库连接失败:", err);
      return;
    }
    return;
  } else if (connection && connection.isConnected) {
    let data: Workload[] = [];
    try {
      // 补充查询必要字段
      data = await connection.query(`
        SELECT workload."id", workload."status", workload."interactionFilter", workload."executeTime"
        from "workload" workload;
      `);
    } catch (err) {
      Logger.error("查询workload数据失败:", err);
      return;
    }

    console.log("jobs", scheduledJobs);

    for (const workload of data) {
      if (scheduledJobs[String(workload.id)]) {
        console.log(`Job with id ${workload.id} exists!`);
        continue;
      }

      Logger.info(`Job with id ${workload.id} does not exist. Creating new scheduled job!`);
      try {
        const filter = await InteractionFilter.findOne({
          where: { id: workload.interactionFilter },
        });

        // 检查filter是否存在,避免后续空指针
        if (!filter) {
          Logger.warn(`未找到ID为${workload.interactionFilter}的InteractionFilter,跳过任务创建`);
          continue;
        }

        Logger.info("FILTER", filter);
        const qb = await connection
          .getRepository(Interaction)
          .createQueryBuilder("interaction")
          .leftJoinAndSelect(
            `interaction.${filter.integrationSystem.toLowerCase()}_metadata`,
            "metadata"
          )
          .where("interaction.system = :system", {
            system: filter.integrationSystem,
          })
          .orderBy("interaction.createdAt", "ASC")
          .getMany();

        const filter_values = await connection
          .getRepository(InteractionFilter_Metadata)
          .createQueryBuilder("filtervalue")
          .leftJoinAndSelect("filtervalue.metadata", "metadata")
          .where("filtervalue.filterId = :filterId", {
            filterId: workload.interactionFilter,
          })
          .distinct()
          .getMany();

        Logger.info("FILTER VALUES", filter_values);

        scheduleJob(
          String(workload.id),
          dayjs(workload.executeTime).toDate(),
          function () {
            console.log("TEST");
          }
        );
        console.log("SCHEDULE TASK AFTER", scheduledJobs);
      } catch (err) {
        Logger.error(`创建workload ${workload.id}对应的任务失败:`, err);
      }
    }
  }
});

expose(() => true);

关键修改点

  • 所有异步操作添加try/catch,捕获异常并记录日志,确保任务流程不中断。
  • 查询workload时补充必要字段,保证后续逻辑有完整数据可用。
  • 增加filter空值检查,避免后续使用filter.integrationSystem时抛出空指针异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 00:07:12