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

如何在Express+GCP环境中实现长时系统迁移任务的实时进度追踪

实现长时系统迁移任务的异步进度追踪方案

基于你的技术栈(Node.js/Express/SSE/GCP),我们可以通过异步任务托管+实时状态存储+SSE推送的组合方案来实现需求,以下是具体实现步骤:

一、整体架构思路

  1. 前端触发迁移后,后端立即返回任务ID,不阻塞用户操作
  2. 迁移任务交由GCP的异步服务(Cloud Tasks/Cloud Function)执行,避免占用Express进程
  3. 任务执行过程中,实时将进度写入GCP的实时数据库(Firestore)
  4. 用户返回页面时,通过SSE接口订阅任务进度,从Firestore获取实时更新并推送给前端

二、后端实现(Express + GCP)

1. 创建迁移任务接口

该接口负责生成任务ID、提交异步任务到GCP,然后立即返回结果:

const express = require('express');
const { v4: uuidv4 } = require('uuid');
const { CloudTasksClient } = require('@google-cloud/tasks');
const router = express.Router();

const tasksClient = new CloudTasksClient();
const PROJECT_ID = 'your-gcp-project-id';
const LOCATION = 'us-central1';
const QUEUE_NAME = 'system-migrate-queue';

// POST /api/system/:systemId/migrate
router.post('/:systemId/migrate', async (req, res) => {
  const { systemId } = req.params;
  const taskId = uuidv4();

  // 构建Cloud Tasks任务请求
  const parent = tasksClient.queuePath(PROJECT_ID, LOCATION, QUEUE_NAME);
  const task = {
    httpRequest: {
      httpMethod: 'POST',
      url: `https://your-cloud-function-url/run-migrate`, // 指向执行迁移的Cloud Function
      body: Buffer.from(JSON.stringify({ systemId, taskId })).toString('base64'),
      headers: {
        'Content-Type': 'application/json',
      },
    },
  };

  // 提交任务到Cloud Tasks
  await tasksClient.createTask({ parent, task });

  // 立即返回任务ID给前端
  res.json({ taskId });
});

2. 进度查询与SSE推送接口

该接口建立SSE连接,先返回当前进度,再实时推送后续更新:

const { Firestore } = require('@google-cloud/firestore');
const firestore = new Firestore();

// GET /api/system/:systemId/migrate/:taskId
router.get('/:systemId/migrate/:taskId', async (req, res) => {
  const { taskId } = req.params;

  // 设置SSE响应头
  res.setHeader('Content-Type', 'text/event-stream');
  res.setHeader('Cache-Control', 'no-cache');
  res.setHeader('Connection', 'keep-alive');
  res.flushHeaders();

  // 获取任务当前状态
  const taskDoc = firestore.collection('migrate-tasks').doc(taskId);
  const initialSnapshot = await taskDoc.get();
  if (!initialSnapshot.exists) {
    res.write(`data: ${JSON.stringify({ status: 'not_found', progress: 0 })}\n\n`);
    res.end();
    return;
  }
  res.write(`data: ${JSON.stringify(initialSnapshot.data())}\n\n`);

  // 监听Firestore文档变化,实时推送进度
  const unsubscribe = taskDoc.onSnapshot(snapshot => {
    const data = snapshot.data();
    res.write(`data: ${JSON.stringify(data)}\n\n`);

    // 任务完成/失败时关闭连接
    if (data.status === 'completed' || data.status === 'failed') {
      unsubscribe();
      res.end();
    }
  });

  // 客户端断开连接时清理监听
  req.on('close', () => {
    unsubscribe();
    res.end();
  });
});

3. GCP Cloud Function:执行迁移任务

该函数负责实际执行迁移逻辑,并实时更新进度到Firestore:

const { Firestore } = require('@google-cloud/firestore');
const firestore = new Firestore();

exports.runMigrate = async (req, res) => {
  const { systemId, taskId } = req.body;
  const taskDoc = firestore.collection('migrate-tasks').doc(taskId);

  // 初始化任务状态
  await taskDoc.set({
    systemId,
    status: 'running',
    progress: 0,
    message: '迁移任务已启动'
  });

  try {
    // 模拟迁移步骤(替换为实际业务逻辑)
    await step1(taskId);
    await step2(taskId);
    await step3(taskId);

    // 任务完成
    await taskDoc.update({
      status: 'completed',
      progress: 100,
      message: '迁移完成'
    });
  } catch (err) {
    // 任务失败
    await taskDoc.update({
      status: 'failed',
      progress: 0,
      message: `迁移失败:${err.message}`
    });
  }

  res.status(200).send('Migration task executed');
};

// 示例迁移步骤,更新进度
async function step1(taskId) {
  const taskDoc = firestore.collection('migrate-tasks').doc(taskId);
  await taskDoc.update({ progress: 30, message: '正在迁移系统配置' });
  await new Promise(resolve => setTimeout(resolve, 10000)); // 模拟耗时操作
}

async function step2(taskId) {
  const taskDoc = firestore.collection('migrate-tasks').doc(taskId);
  await taskDoc.update({ progress: 60, message: '正在迁移用户数据' });
  await new Promise(resolve => setTimeout(resolve, 15000));
}

async function step3(taskId) {
  const taskDoc = firestore.collection('migrate-tasks').doc(taskId);
  await taskDoc.update({ progress: 90, message: '正在验证迁移结果' });
  await new Promise(resolve => setTimeout(resolve, 5000));
}

三、前端实现

1. 触发迁移并存储任务ID

// 点击「迁移系统」按钮时的逻辑
async function triggerMigration(systemId) {
  try {
    const res = await fetch(`/api/system/${systemId}/migrate`, {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' }
    });
    const { taskId } = await res.json();
    // 将任务ID与systemId关联存储到localStorage
    localStorage.setItem(`migrate-task-${systemId}`, taskId);
    alert('迁移任务已启动,可稍后返回页面查看进度');
  } catch (err) {
    console.error('触发迁移失败:', err);
  }
}

2. 页面返回后监听进度

// 进入/system/:systemId页面时的逻辑
async function initProgressMonitor(systemId) {
  const taskId = localStorage.getItem(`migrate-task-${systemId}`);
  if (!taskId) return;

  // 建立SSE连接
  const eventSource = new EventSource(`/api/system/${systemId}/migrate/${taskId}`);

  eventSource.onmessage = (event) => {
    const data = JSON.parse(event.data);
    // 更新页面UI
    document.getElementById('progress-bar').style.width = `${data.progress}%`;
    document.getElementById('status-text').textContent = data.message;

    // 任务完成/失败时关闭连接并清理存储
    if (data.status === 'completed' || data.status === 'failed') {
      eventSource.close();
      localStorage.removeItem(`migrate-task-${systemId}`);
    }
  };

  eventSource.onerror = (err) => {
    console.error('SSE连接出错:', err);
    eventSource.close();
  };
}

四、GCP配置要点

  1. Cloud Tasks队列:创建一个队列,配置合适的重试策略(比如失败后重试3次),确保任务可靠执行
  2. Firestore集合:创建migrate-tasks集合,用于存储任务状态,无需预定义结构,动态写入即可
  3. 权限配置:
    • 给Express服务账号分配Cloud Tasks Enqueuer权限,允许提交任务
    • 给Cloud Function服务账号分配Firestore Editor权限,允许读写任务状态
  4. 超时设置:Cloud Function的超时时间设置为60分钟(最大值),适配最长1小时的迁移任务

内容的提问来源于stack exchange,提问作者G Madhabananda Patra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 03:10:02