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

Chrome扩展流数据缓冲区并发读写安全处理及替代方案咨询

流式数据缓冲与并发安全处理方案

一、基于数组的安全实现(单线程环境无竞态)

JavaScript是单线程事件循环模型,不存在真正的并发读写冲突,但需通过状态标记区分缓冲/转发阶段,避免数据重复或丢失:

let buffer = [];
let isStreaming = false;

// 处理Socket.io流式数据
socket.on('stream', (wordChunk) => {
  if (isStreaming) {
    // 已触发推送,直接发送至内容脚本
    port.postMessage({ msg: 'streamData', wordChunk });
  } else {
    // 未触发,存入缓冲区
    buffer.push(wordChunk);
  }
});

// 处理内容脚本的startStreaming消息
if (msg.msg === 'startStreaming') {
  isStreaming = true;
  // 发送缓冲的历史数据
  buffer.forEach(chunk => {
    port.postMessage({ msg: 'streamData', chunk });
  });
  // 清空缓冲区,防止重复发送旧数据
  buffer = [];
}

核心要点:

  • isStreaming状态切换为原子操作,单线程下不会出现并发问题
  • 遍历缓冲区时,Socket回调需等待当前同步代码执行完毕才会触发,保证历史数据完整发送
  • 发送完成后清空缓冲区,避免用户多次触发startStreaming时重复推送

针对多次悬停预加载场景,可在每次悬停时重置状态:

function onHover() {
  // 清空旧缓冲,重置状态
  buffer = [];
  isStreaming = false;
  // 重新订阅Socket流(若需重新拉取最新数据)
  socket.off('stream');
  socket.on('stream', (wordChunk) => {
    isStreaming ? port.postMessage({ msg: 'streamData', wordChunk }) : buffer.push(wordChunk);
  });
}

二、替代数组的高效缓冲区方案

若流式数据量较大,数组头部操作(如shift)性能较低,可采用以下方案:

1. 双数组实现高效队列

用inbox和outbox模拟队列,所有操作均为均摊O(1)时间复杂度:

class EfficientQueue {
  constructor() {
    this.inbox = [];
    this.outbox = [];
  }

  enqueue(item) {
    this.inbox.push(item);
  }

  dequeue() {
    if (this.outbox.length === 0) {
      this.outbox = this.inbox.reverse();
      this.inbox = [];
    }
    return this.outbox.pop();
  }

  getAll() {
    const allItems = [...this.outbox.reverse(), ...this.inbox];
    this.inbox = [];
    this.outbox = [];
    return allItems;
  }

  isEmpty() {
    return this.inbox.length === 0 && this.outbox.length === 0;
  }
}

// 使用示例
const buffer = new EfficientQueue();
let isStreaming = false;

socket.on('stream', (chunk) => {
  isStreaming ? port.postMessage({ msg: 'streamData', chunk }) : buffer.enqueue(chunk);
});

if (msg.msg === 'startStreaming') {
  isStreaming = true;
  buffer.getAll().forEach(chunk => port.postMessage({ msg: 'streamData', chunk }));
}

2. 使用原生ReadableStream API

原生ReadableStream支持自动缓冲、背压处理(内容脚本处理缓慢时自动暂停接收),适合大流量场景:

let streamController;
let isStreamStarted = false;

// 创建可控制的可读流
const readableStream = new ReadableStream({
  start(controller) {
    streamController = controller;
  }
});

// 处理Socket数据
socket.on('stream', (chunk) => {
  streamController.enqueue(new TextEncoder().encode(chunk));
});

// 处理startStreaming消息
if (msg.msg === 'startStreaming') {
  isStreamStarted = true;
  const reader = readableStream.getReader();
  async function forwardStream() {
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      port.postMessage({ msg: 'streamData', chunk: new TextDecoder().decode(value) });
    }
  }
  forwardStream();
}

三、多次悬停场景的最佳实践

  • 每次悬停重置缓冲:清空旧数据,重新接收最新流式内容,避免点击时展示过期数据
  • 取消旧订阅:每次悬停时移除之前的Socketstream事件监听,防止旧数据混入新缓冲
  • 设置缓冲超时:若用户悬停后长时间未点击,自动清空缓冲,减少内存占用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 21:22:49