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

React与Nest.js应用中WebSocket无法持续发送消息的排查求助

问题描述

我需要从React客户端向Nest.js服务器发送数据,服务器执行耗时的复杂计算后,通过WebSocket将计算进度实时回传给客户端。但现在遇到问题:服务器在循环里调用client.emit()发送进度消息时,客户端只能收到并打印第一条消息。

客户端代码

const socket = io('http://localhost:5000');    

socket.emit('correlate', {
   formData: data
}, (res) => {
     // 一些副作用操作            
});

// 监听进度
socket.on('correlate', (d) => {
   console.log(d); // 这里只打印第一条消息
});

服务器代码

for(const relationRow of relationSheet) {
   // 复杂计算逻辑

   client.emit('correlate', 70 * (i / relationSheet.length / 100));
}
解决方案

一、修复WebSocket消息不实时的问题

1. 解决事件名称冲突

你客户端和服务器都用了correlate这个事件名,客户端既用它发送请求,又用它监听进度,直接导致了事件冲突。把进度监听的事件名改成不一样的,比如correlate-progress:

客户端修改后代码:

// 监听进度
socket.on('correlate-progress', (d) => {
   console.log(d); // 现在能接收所有进度消息
});

服务器修改后代码:

// 注意要通过entries()获取索引i,原代码里的i未定义
for(const [i, relationRow] of relationSheet.entries()) { 
   // 复杂计算逻辑

   client.emit('correlate-progress', 70 * (i / relationSheet.length / 100));
}

2. 避免阻塞事件循环

如果你的复杂计算是同步逻辑,会直接阻塞Node.js的事件循环,导致WebSocket消息无法及时推送出去。可以把计算拆成异步任务,用setImmediate或者Promise让事件循环有间隙处理消息:

服务器修改后代码:

async function processRows() {
  for(const [i, relationRow] of relationSheet.entries()) {
    // 把同步计算包进Promise,让事件循环能喘息
    await new Promise(resolve => {
      // 复杂计算逻辑
      resolve();
    });
    // 用setImmediate确保消息被及时发送
    setImmediate(() => {
      client.emit('correlate-progress', 70 * (i / relationSheet.length / 100));
    });
  }
}

processRows();

二、无需WebSocket的替代方案

1. 服务器发送事件(SSE)

SSE是单向的服务器向客户端推送技术,适合进度推送场景,比WebSocket实现更简单,不需要维护双向连接:

  • React客户端代码:
function useProgressSSE(taskId) {
  useEffect(() => {
    const eventSource = new EventSource(`http://localhost:5000/correlate-progress/${taskId}`);
    eventSource.onmessage = (event) => {
      const progress = JSON.parse(event.data);
      console.log(progress);
    };
    eventSource.onerror = () => {
      eventSource.close();
    };
    return () => {
      eventSource.close();
    };
  }, [taskId]);
}

// 先发送计算请求获取任务ID
fetch('http://localhost:5000/correlate', {
  method: 'POST',
  body: JSON.stringify({ formData: data }),
  headers: { 'Content-Type': 'application/json' }
})
.then(res => res.json())
.then(({ taskId }) => {
  useProgressSSE(taskId);
});
  • Nest.js服务器代码:
import { Controller, Post, Body, Res, Get, Param } from '@nestjs/common';
import { Response } from 'express';
import { v4 as uuidv4 } from 'uuid';

@Controller('correlate')
export class CorrelateController {
  private taskProgress = new Map<string, number>();

  @Post()
  async startCorrelate(@Body() body: any, @Res() res: Response) {
    const taskId = uuidv4();
    this.taskProgress.set(taskId, 0);
    // 异步启动计算任务
    this.processTask(taskId, body.formData);
    res.json({ taskId });
  }

  private async processTask(taskId: string, formData: any) {
    const relationSheet = await this.getRelationSheet(formData);
    for(const [i, row] of relationSheet.entries()) {
      // 复杂计算逻辑
      const progress = 70 * (i / relationSheet.length / 100);
      this.taskProgress.set(taskId, progress);
    }
    this.taskProgress.set(taskId, 100);
  }

  @Get('progress/:taskId')
  async getProgress(@Param('taskId') taskId: string, @Res() res: Response) {
    res.setHeader('Content-Type', 'text/event-stream');
    res.setHeader('Cache-Control', 'no-cache');
    res.setHeader('Connection', 'keep-alive');

    const interval = setInterval(() => {
      const progress = this.taskProgress.get(taskId) || 0;
      res.write(`data: ${JSON.stringify(progress)}\n\n`);
      if (progress >= 100) {
        clearInterval(interval);
        res.end();
      }
    }, 100);

    res.on('close', () => {
      clearInterval(interval);
    });
  }

  private async getRelationSheet(formData: any) {
    // 模拟获取数据逻辑
    return new Promise(resolve => setTimeout(() => resolve(Array(100).fill({})), 100));
  }
}

2. 轮询(Polling)

客户端定期向服务器请求当前进度,实现最简单但效率较低,适合对实时性要求不高的场景:

  • React客户端代码:
// 发送计算请求
fetch('http://localhost:5000/correlate', {
  method: 'POST',
  body: JSON.stringify({ formData: data }),
  headers: { 'Content-Type': 'application/json' }
})
.then(res => res.json())
.then(({ taskId }) => {
  // 每500ms轮询一次进度
  const pollInterval = setInterval(() => {
    fetch(`http://localhost:5000/correlate-progress/${taskId}`)
    .then(res => res.json())
    .then(progress => {
      console.log(progress);
      if (progress >= 100) {
        clearInterval(pollInterval);
      }
    });
  }, 500);
});
  • Nest.js服务器代码:
import { Controller, Post, Body, Get, Param } from '@nestjs/common';
import { v4 as uuidv4 } from 'uuid';

@Controller('correlate')
export class CorrelateController {
  private taskProgress = new Map<string, number>();

  @Post()
  async startCorrelate(@Body() body: any) {
    const taskId = uuidv4();
    this.taskProgress.set(taskId, 0);
    // 异步执行计算
    (async () => {
      const relationSheet = await this.getRelationSheet(body.formData);
      for(const [i, row] of relationSheet.entries()) {
        // 复杂计算逻辑
        const progress = 70 * (i / relationSheet.length / 100);
        this.taskProgress.set(taskId, progress);
      }
      this.taskProgress.set(taskId, 100);
    })();
    return { taskId };
  }

  @Get('progress/:taskId')
  getProgress(@Param('taskId') taskId: string) {
    return this.taskProgress.get(taskId) || 0;
  }

  private async getRelationSheet(formData: any) {
    // 模拟获取数据逻辑
    return new Promise(resolve => setTimeout(() => resolve(Array(100).fill({})), 100));
  }
}

内容的提问来源于stack exchange,提问作者Sa1m0n

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 08:32:44