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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 23:31:00