如何让Angular客户端等待Node后端返回Airflow任务完成通知?
实现客户端等待Airflow任务完成后响应的方案
针对你描述的场景,有三种可行的方案让Angular客户端等到Airflow任务完成后再拿到后端响应:
1. 长轮询(Long Polling)
这是最贴合需求的方案,核心是让后端保持请求连接直到任务完成:
- Angular客户端向Node后端发起请求,携带任务标识相关参数
- Node后端收到请求后,触发Airflow任务,然后不立即返回响应,挂起该请求
- 当Airflow任务完成后,通过HTTP回调通知Node后端任务结果
- Node后端收到回调后,立刻把结果返回给挂起的客户端请求
- Angular客户端拿到响应后结束等待
代码示例
Node后端(Express)
const uuid = require('uuid'); const pendingRequests = {}; // 处理客户端触发请求的路由 app.post('/trigger-airflow', async (req, res) => { const taskId = uuid.v4(); // 触发Airflow任务,传入任务ID和回调地址 await triggerAirflowTask(taskId, `${process.env.BASE_URL}/airflow-callback/${taskId}`); // 保存当前响应对象,等待回调触发 pendingRequests[taskId] = res; }); // Airflow任务完成后的回调路由 app.post('/airflow-callback/:taskId', (req, res) => { const { taskId } = req.params; const taskResult = req.body; // 找到对应挂起的请求并返回结果 if (pendingRequests[taskId]) { pendingRequests[taskId].json({ success: true, result: taskResult }); delete pendingRequests[taskId]; } res.sendStatus(200); });
Angular客户端
import { HttpClient } from '@angular/common/http'; import { Injectable } from '@angular/core'; @Injectable() export class AirflowService { constructor(private http: HttpClient) {} async triggerAndWaitForResult() { try { const response = await this.http.post('/trigger-airflow', {}, { observe: 'response' }).toPromise(); console.log('任务完成结果:', response.body); } catch (error) { console.error('请求出错:', error); } } }
2. WebSocket 实时通信
通过WebSocket建立客户端和后端的持久连接,任务完成后后端主动推送结果:
- Angular客户端和Node后端建立WebSocket连接
- 客户端发送消息触发Airflow任务,携带任务ID
- Node后端触发Airflow任务,记录任务ID对应的WebSocket连接
- Airflow任务完成后回调后端,后端通过对应的WebSocket连接把结果推送给客户端
- 客户端收到推送后处理结果,可选择关闭连接
代码示例
Node后端(WebSocket)
const WebSocket = require('ws'); const uuid = require('uuid'); const wss = new WebSocket.Server({ port: 8080 }); const taskConnections = {}; wss.on('connection', (ws) => { ws.on('message', async (message) => { const data = JSON.parse(message); if (data.type === 'trigger') { const taskId = uuid.v4(); taskConnections[taskId] = ws; // 触发Airflow任务 await triggerAirflowTask(taskId, `${process.env.BASE_URL}/airflow-callback/${taskId}`); } }); }); // Airflow回调路由 app.post('/airflow-callback/:taskId', (req, res) => { const { taskId } = req.params; const taskResult = req.body; if (taskConnections[taskId]) { taskConnections[taskId].send(JSON.stringify({ success: true, result: taskResult })); delete taskConnections[taskId]; } res.sendStatus(200); });
Angular客户端
import { webSocket } from 'rxjs/webSocket'; triggerTask() { const socket$ = webSocket('ws://localhost:8080'); socket$.subscribe({ next: (msg) => { console.log('任务完成结果:', msg); socket$.complete(); }, error: (err) => console.error('Socket错误:', err), complete: () => console.log('Socket连接关闭') }); socket$.next({ type: 'trigger' }); }
3. 客户端定时轮询
这是最简单的方案,客户端主动定期查询任务状态:
- Angular客户端发起请求触发Airflow任务,后端返回任务ID
- 客户端每隔固定时间(比如5秒)用任务ID向后端查询任务状态
- 后端查询Airflow任务状态(或通过回调记录的状态),返回给客户端
- 当客户端收到任务完成的状态时,停止轮询并处理结果
代码示例
Angular客户端
async triggerTaskAndPoll() { // 触发任务,获取任务ID const triggerRes = await this.http.post('/trigger-airflow', {}).toPromise(); const taskId = triggerRes['taskId']; // 定时轮询任务状态 const pollInterval = setInterval(async () => { const statusRes = await this.http.get(`/task-status/${taskId}`).toPromise(); if (statusRes['status'] === 'completed') { clearInterval(pollInterval); console.log('任务完成结果:', statusRes['result']); } else if (statusRes['status'] === 'failed') { clearInterval(pollInterval); console.error('任务失败'); } }, 5000); // 每5秒查询一次 }
Node后端
const uuid = require('uuid'); const taskStatus = {}; app.post('/trigger-airflow', async (req, res) => { const taskId = uuid.v4(); taskStatus[taskId] = { status: 'running' }; await triggerAirflowTask(taskId, `${process.env.BASE_URL}/airflow-callback/${taskId}`); res.json({ taskId }); }); app.get('/task-status/:taskId', (req, res) => { const { taskId } = req.params; res.json(taskStatus[taskId] || { status: 'not_found' }); }); app.post('/airflow-callback/:taskId', (req, res) => { const { taskId } = req.params; taskStatus[taskId] = { status: 'completed', result: req.body }; res.sendStatus(200); });
内容的提问来源于stack exchange,提问作者Splinteer
相关产品推荐
相关产品推荐

