22 最终综合项目:外部分析任务服务

本项目把 Express、数据库、后台 Worker 和 child_process.spawn() 连接起来。文档提供契约和骨架,但不直接给出完整进程管理答案。

1. 业务场景

管理员提交一个已经上传的数据文件,选择受支持的分析程序:

summary       Node.js 测试脚本
python-stats  Python 统计脚本
r-report      Rscript 报表脚本

服务不能让用户提交任意 command 或脚本路径,只能从服务端白名单映射:

type AnalysisType =
  | "summary"
  | "python-stats"
  | "r-report";

2. API 契约

创建:

POST /api/analysis-tasks
Authorization: Bearer <token>
Content-Type: application/json

{
  "analysisType": "python-stats",
  "inputFileId": "812",
  "parameters": {
    "method": "mean"
  }
}

成功:

{
  "code": "OK",
  "message": "任务已创建",
  "data": {
    "taskId": "1092",
    "status": "queued"
  }
}

其他接口:

GET  /api/analysis-tasks/:taskId
POST /api/analysis-tasks/:taskId/cancel
GET  /api/analysis-tasks/:taskId/logs
GET  /api/analysis-tasks/:taskId/download

3. 状态机

queued
  → starting
  → running
  → succeeded

queued/running → cancelling → cancelled
starting/running → failed
running → timed_out

状态不能只放在当前 Node.js 进程的 Map 中。PM2/Cluster 重启和多进程会让内存状态丢失或不一致。

4. 数据类型骨架

export interface AnalysisTaskPayload {
  taskId: string;
  analysisType: AnalysisType;
  inputPath: string;
  outputDirectory: string;
  parameters: Readonly<Record<string, string | number | boolean>>;
}

export interface RunningProcessInfo {
  taskId: string;
  childPid: number;
  startedAt: Date;
  abortController: AbortController;
}

export type AnalysisProgressEvent = {
  type: "progress";
  percent: number;
  message?: string;
};

export type AnalysisResultEvent = {
  type: "result";
  outputFileName: string;
  metrics: Record<string, number>;
};

message? 在启用 exactOptionalPropertyTypes 时,构造对象不能随意赋 undefined,需要省略字段或把类型显式写成 string | undefined

5. 命令白名单

export interface ProcessCommand {
  executable: string;
  args: string[];
  cwd: string;
  env: NodeJS.ProcessEnv;
}

export function buildProcessCommand(
  task: AnalysisTaskPayload,
): ProcessCommand {
  switch (task.analysisType) {
    case "python-stats":
      // TODO:从可信配置取得 python3 和脚本绝对路径
      // TODO:按数组构造参数,不使用 shell: true
      throw new Error("TODO");
    case "r-report":
      throw new Error("TODO");
    case "summary":
      throw new Error("TODO");
  }
}

输入文件必须先解析到允许的 Storage 根目录并验证任务有权访问。参数数组能降低 Shell 注入风险,但不能代替路径授权和业务参数校验。

6. 子进程输出协议

脚本 stdout 使用 JSON Lines,一行一个事件:

{"type":"progress","percent":10,"message":"读取输入"}
{"type":"progress","percent":80,"message":"生成结果"}
{"type":"result","outputFileName":"result.json","metrics":{"rows":120}}

stderr 只用于诊断和进度文本,不作为最终机器结果。

退出约定:

0  程序成功,并且必须存在合法 result 事件和输出文件
2  输入参数错误
3  输入数据错误
4  外部依赖失败
5  内部分析失败

即使退出码为 0,如果没有合法结果事件、输出路径越界或文件不存在,任务仍应判定失败。

7. 运行骨架

export async function runAnalysisProcess(
  task: AnalysisTaskPayload,
  signal: AbortSignal,
): Promise<AnalysisResultEvent> {
  const command = buildProcessCommand(task);

  // TODO 1:spawn executable + args,明确 cwd/env/stdio/signal
  // TODO 2:按行解析 stdout,限制行长和总字节
  // TODO 3:有界保存 stderr,不把 Chunk 当成行
  // TODO 4:同时处理 error 和 close,防止重复 settle
  // TODO 5:非零退出码映射业务错误
  // TODO 6:验证 result 事件和输出文件
  // TODO 7:失败、超时、取消后清理半成品

  throw new Error("TODO");
}

8. 任务领取与并发

单进程时简单轮询看似可行;Cluster/PM2 多实例后,两个 Worker 可能同时领取同一任务。必须使用一种原子领取机制:

数据库条件 UPDATE
SELECT ... FOR UPDATE SKIP LOCKED(结合版本与事务验证)
BullMQ/Redis 队列
专门的任务调度系统

任务级并发必须有上限。每个服务进程各允许 4 个任务,在 4 个实例上就是最多 16 个外部子进程。

9. 取消与超时

取消接口不能仅修改数据库状态:

前端请求取消
  → 原子标记 cancelling
  → 通知拥有该任务的 Worker
  → AbortController.abort()
  → 子进程收到终止请求
  → 等待 close
  → 清理半成品
  → 标记 cancelled

如果 Worker 已崩溃,需要恢复程序根据 PID、租约或任务心跳判断任务是否遗留。不能保存 PID 后盲目 kill,因为 PID 可能被操作系统复用;还要校验 Worker 所有权和任务实例标识。

超时也走相同的受控终止流程。child.kill() 返回成功不等于退出,仍需等待 close

10. 优雅停机

Worker 收到 SIGTERM:

停止领取新任务
  → 将运行任务设为 draining/保持租约
  → 按策略等待或取消子进程
  → 等待 close
  → 关闭数据库/Redis
  → 设置 exitCode 并自然退出

必须设置最大宽限期。Primary/PM2/Kubernetes 的超时应大于应用内部宽限期,避免应用还在清理时被强制杀死。

11. 输出文件安全

12. 错误码建议

ANALYSIS_TYPE_UNSUPPORTED
ANALYSIS_TASK_NOT_FOUND
ANALYSIS_TASK_NOT_READY
ANALYSIS_TASK_ALREADY_FINISHED
ANALYSIS_PROCESS_START_FAILED
ANALYSIS_PROCESS_FAILED
ANALYSIS_PROCESS_TIMEOUT
ANALYSIS_PROCESS_CANCELLED
ANALYSIS_OUTPUT_INVALID
ANALYSIS_OUTPUT_MISSING
ANALYSIS_WORKER_UNAVAILABLE

未知程序错误不能伪装成“输入错误”。日志记录原始 cause,公开响应使用稳定消息。

13. 必测场景

命令不存在
脚本路径不存在
cwd 不存在
输入文件无权限
stdout 合法 JSONL
stdout 出现非法 JSON
JSON 行跨多个 Chunk
stdout 超过上限
stderr 有内容但 exit 0
exit 非0但 stderr 为空
exit 0但缺少结果事件
结果路径越过 outputDirectory
执行超时
用户主动取消
Signal 后脚本延迟退出
脚本创建孙进程
Worker 执行中崩溃
两个 Worker 同时尝试领取任务
PM2 reload

14. 验收清单

[ ] 用户不能决定任意 executable/script
[ ] command 与 args 分离,未开启 shell:true
[ ] cwd 和 env 明确且受限
[ ] stdout/stderr 始终被消费或明确 inherit/ignore
[ ] Chunk 不被当作完整行
[ ] 同时处理启动 error 和最终 close
[ ] 退出码、Signal、结果协议共同决定成功
[ ] 超时/取消后等待实际 close
[ ] 子进程并发有全局或队列级上限
[ ] Cluster 多实例不会重复领取任务
[ ] SIGTERM 有优雅停机和最大宽限期
[ ] 日志不含 Token、密码和敏感样本信息

练习题

  1. 实现 JSON Lines 解析器,正确保留跨 Chunk 的半行。
  2. 设计一个只保留 stderr 最后 64 KiB 的环形缓冲。
  3. 用数据库条件更新设计任务原子领取伪代码。
  4. 解释为什么退出码 0 仍不能单独证明分析结果有效。
  5. 设计 Worker 崩溃后的租约和任务恢复策略。