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

在Express中使用Server-Sent Events实现音频处理进度反馈

如何在POST音频处理流程中向SSE连接发送实时进度更新?

问题背景

我正在开发一个音频处理应用,流程是用户从浏览器上传音频文件,服务器通过POST /api/track接收文件后执行8个耗时的处理步骤,完成后返回结果。现在想在每个步骤完成时向浏览器发送增量反馈(比如“文件已上传”“元数据已处理”等),已经确定用Server-Sent Events(SSE)实现,但遇到了一个卡点:

我的POST处理逻辑在audio.process函数中,而SSE的GET路由/stream是单独的端点。目前测试SSE是正常的,但不知道怎么从audio.process里把步骤消息发送到对应的SSE连接中。我在想是不是要在audio.process里发起GET请求到/stream,但感觉这不是最优方案,想寻求更好的实现思路。

现有代码片段:

服务器端POST路由

app.post('/api/track', upload.single('track'), audio.process)

服务器端测试SSE路由

app.get('/stream', function(req, res) {
  res.sseSetup()
  for (var i = 0; i < 5; i++) {
    res.sseSend({count: i})
  }
})

客户端SSE监听代码

progress : () => {
  if (!!window.EventSource) {
    const source = new EventSource('/stream')
    source.addEventListener('message', function(e) {
      let data = JSON.parse(e.data)
      console.log(e);
    }, false)
    source.addEventListener('open', function(e) {
      console.log("Connected to /stream");
    }, false)
    source.addEventListener('error', function(e) {
      if (e.target.readyState == EventSource.CLOSED) {
        console.log("Disconnected from /stream");
      } else if (e.target.readyState == EventSource.CONNECTING) {
        console.log('Connecting to /stream');
      }
    }, false)
  } else {
    console.log("Your browser doesn't support SSE")
  }
}

解决方案思路

首先明确:POST请求和SSE的GET请求是两个完全独立的HTTP连接,不能直接在audio.process里调用SSE路由的方法。我们需要一个共享的通信中间层来连接这两个请求上下文,下面是几种可行的方案,按推荐优先级排序:

1. 使用Node.js的EventEmitter作为事件总线(单进程场景首选)

这是最直接且轻量的方案,在服务器端创建一个全局的事件发射器,让SSE连接订阅特定事件,而audio.process在每个步骤完成时发布对应事件。

具体实现步骤:

  • 首先,在服务器初始化时创建一个全局的EventEmitter实例:
const EventEmitter = require('events');
const progressEmitter = new EventEmitter();
// 可选:设置最大监听器数,避免多连接时的警告
progressEmitter.setMaxListeners(0);
  • 修改SSE路由,让每个连接绑定到唯一的任务ID(需要客户端传递任务ID来区分不同用户的进度):
app.get('/stream', function(req, res) {
  // 配置SSE响应头
  res.writeHead(200, {
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'Connection': 'keep-alive'
  });
  res.flushHeaders();

  // 从请求参数获取任务ID(客户端打开SSE时需要传递)
  const taskId = req.query.taskId;

  // 定义消息发送函数
  const sendProgress = (message) => {
    res.write(`data: ${JSON.stringify({ message })}\n\n`);
  };

  // 订阅当前任务的进度事件
  progressEmitter.on(`progress:${taskId}`, sendProgress);

  // 连接关闭时取消订阅,防止内存泄漏
  req.on('close', () => {
    progressEmitter.off(`progress:${taskId}`, sendProgress);
  });
})
  • 调整客户端逻辑:上传文件前生成唯一taskId,先打开SSE连接监听该任务,再携带taskId发起POST上传:
// 生成唯一任务ID(可以用更严谨的生成方式,比如uuid)
const taskId = Date.now().toString() + Math.random().toString(36).slice(2, 10);

// 先启动SSE连接,监听当前任务的进度
const source = new EventSource(`/stream?taskId=${taskId}`);
source.addEventListener('message', (e) => {
  const data = JSON.parse(e.data);
  console.log('进度更新:', data.message);
  // 这里可以把消息渲染到页面上,比如更新进度条或提示文本
});

// 构造表单数据,携带taskId和音频文件
const formData = new FormData();
formData.append('track', audioFile);
formData.append('taskId', taskId);

// 发起POST上传请求
fetch('/api/track', {
  method: 'POST',
  body: formData
})
.then(res => res.json())
.then(result => {
  console.log('处理完成:', result);
  // 处理完成后主动关闭SSE连接
  source.close();
})
.catch(err => {
  console.error('上传失败:', err);
  source.close();
});
  • 修改audio.process函数,在每个步骤完成时通过事件发射器发送消息:
async function process(req, res) {
  const taskId = req.body.taskId;
  try {
    // 步骤1:文件已上传,发送第一个进度消息
    progressEmitter.emit(`progress:${taskId}`, '文件已上传');
    
    // 步骤2:处理元数据
    await processMetadata(req.file);
    progressEmitter.emit(`progress:${taskId}`, '元数据已处理');
    
    // 步骤3:提取封面图片
    await extractCoverImage(req.file);
    progressEmitter.emit(`progress:${taskId}`, '图片已提取');
    
    // ... 依次处理剩余步骤,每个步骤完成后发送对应消息
    
    // 所有步骤完成,返回最终结果
    res.json({ status: 'success', message: '全部处理完成' });
  } catch (err) {
    // 处理出错时发送错误消息
    progressEmitter.emit(`progress:${taskId}`, `处理失败: ${err.message}`);
    res.status(500).json({ status: 'error', message: err.message });
  }
}

2. 使用Redis Pub/Sub(多进程/分布式场景首选)

如果你的应用是多进程部署或者分布式架构,内存中的EventEmitter无法跨进程通信,这时候可以用Redis的发布/订阅功能来替代EventEmitter:

  • SSE连接启动时,订阅Redis的progress:{taskId}频道
  • audio.process在每个步骤完成时,向progress:{taskId}频道发布消息
  • Redis会自动把消息推送给所有订阅该频道的SSE连接

这种方案的好处是天然支持跨进程、跨服务器的通信,适合大型应用场景。

3. 绝对不要在audio.process中发起GET请求到/stream

你之前考虑的这个思路是不可行的——每个GET请求会创建一个全新的SSE连接,而不是向已有的客户端连接发送消息。这不仅无法实现进度通知,还会造成服务器资源的浪费,一定要避免。


关键注意事项

  • 任务ID的唯一性:必须用唯一的taskId区分不同用户的任务,否则会出现消息发错人的情况,推荐用uuid库生成更可靠的ID。
  • 内存泄漏防范:SSE连接关闭时一定要取消事件订阅,否则EventEmitter会保留无效的监听器,导致内存泄漏。
  • 心跳机制:有些服务器或代理会自动关闭长时间空闲的连接,建议每30秒发送一个空的心跳消息(res.write('data: \n\n'))来保持连接活跃。
  • 错误处理:处理过程中发生错误时,要及时发送错误消息给客户端,并让客户端关闭SSE连接,避免无效等待。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:44:06