如何在Express+GCP环境中实现长时系统迁移任务的实时进度追踪
实现长时系统迁移任务的异步进度追踪方案
基于你的技术栈(Node.js/Express/SSE/GCP),我们可以通过异步任务托管+实时状态存储+SSE推送的组合方案来实现需求,以下是具体实现步骤:
一、整体架构思路
- 前端触发迁移后,后端立即返回任务ID,不阻塞用户操作
- 迁移任务交由GCP的异步服务(Cloud Tasks/Cloud Function)执行,避免占用Express进程
- 任务执行过程中,实时将进度写入GCP的实时数据库(Firestore)
- 用户返回页面时,通过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配置要点
- Cloud Tasks队列:创建一个队列,配置合适的重试策略(比如失败后重试3次),确保任务可靠执行
- Firestore集合:创建
migrate-tasks集合,用于存储任务状态,无需预定义结构,动态写入即可 - 权限配置:
- 给Express服务账号分配
Cloud Tasks Enqueuer权限,允许提交任务 - 给Cloud Function服务账号分配
Firestore Editor权限,允许读写任务状态
- 给Express服务账号分配
- 超时设置:Cloud Function的超时时间设置为60分钟(最大值),适配最长1小时的迁移任务
内容的提问来源于stack exchange,提问作者G Madhabananda Patra
相关产品推荐
相关产品推荐

