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

React中RxJS WebSocket接收消息时状态不更新问题求助

问题解决:React中RxJS WebSocket流式消息拼接时状态不更新的问题

问题根源

你的useEffect依赖数组是空的,这意味着subscribe回调只会在组件挂载时创建一次,它捕获的是初始状态的currentMessage和messages。后续状态更新后,回调里的变量还是初始值(currentMessage一直是null),所以永远只会进入else分支,无法触发中间的拼接逻辑。

修复方案

方案1:使用函数式更新(推荐)

React的setState支持传入函数,函数会接收当前最新的状态值,这样就能避开闭包捕获旧值的问题。修改subscribe里的逻辑:

import { useContext, useEffect, useState } from "react";
import { Context } from "../context/Context";

function Home() {
  const { subject, getWelcomeMessage } = useContext(Context);
  const [messages, setMessages] = useState([]);
  const [currentMessage, setCurrentMessage] = useState(null);
  const [previousMessage, setPreviousMessage] = useState(null);

  useEffect(() => {
    console.log("currentMessage", currentMessage);
    console.log("messages", messages);
  }, [currentMessage, messages]);

  useEffect(() => {
    const subscription = subject.subscribe({
      next: (val) => {
        console.log("Val", val);
        if (val.m === "$$$EOM$$$" || val.m === "$$$EOM:CRITICAL$$$" || val.m === "$$$EOM:ERROR$$$") {
          console.log("&& 结束标记");
          setCurrentMessage(prev => {
            const updatedMsg = prev ? {...prev, inProgress: false} : null;
            setPreviousMessage(updatedMsg);
            return null;
          });
        } else {
          setCurrentMessage(prev => {
            if (prev !== null) {
              console.log("当前消息非空,继续拼接");
              return {
                ...prev,
                inProgress: true,
                content: [...prev.content, val.m]
              };
            } else {
              console.log("当前消息为空,创建新消息");
              const newMsg = { role: "assistant", content: [val.m] };
              setMessages(prevMsgs => [...prevMsgs, newMsg]);
              return newMsg;
            }
          });
        }
      },
    });
    getWelcomeMessage();

    // 组件卸载时取消订阅,防止内存泄漏
    return () => subscription.unsubscribe();
  }, [subject, getWelcomeMessage]);

  return (
    <div>
      <h3>Batty Chat Bot</h3>
    </div>
  );
}

Home.propTypes = {};

export default Home;

方案2:更新useEffect依赖数组

把currentMessage和messages加入依赖数组,同时用RxJS的takeUntil防止重复订阅:

import { useContext, useEffect, useState, useCallback } from "react";
import { Context } from "../context/Context";
import { Subject } from "rxjs";
import { takeUntil } from "rxjs/operators";

function Home() {
  const { subject, getWelcomeMessage } = useContext(Context);
  const [messages, setMessages] = useState([]);
  const [currentMessage, setCurrentMessage] = useState(null);
  const [previousMessage, setPreviousMessage] = useState(null);
  const unsubscribe$ = new Subject();

  const handleMessage = useCallback((val) => {
    console.log("Val", val);
    if (val.m === "$$$EOM$$$" || val.m === "$$$EOM:CRITICAL$$$" || val.m === "$$$EOM:ERROR$$$") {
      console.log("&& 结束标记");
      const tempCurrentMessage = { ...currentMessage };
      if (tempCurrentMessage) {
        tempCurrentMessage.inProgress = false;
      }
      setPreviousMessage(tempCurrentMessage);
      setCurrentMessage(null);
    } else if (currentMessage !== null) {
      console.log("当前消息非空");
      const tempCurrentMessage = { ...currentMessage };
      tempCurrentMessage.inProgress = true;
      tempCurrentMessage.content = [...currentMessage.content, val.m];
      setCurrentMessage(tempCurrentMessage);
    } else {
      console.log("当前消息为空");
      const tempCurrentMessage = { role: "assistant", content: [val.m] };
      setCurrentMessage(tempCurrentMessage);
      setMessages([...messages, tempCurrentMessage]);
    }
  }, [currentMessage, messages]);

  useEffect(() => {
    subject.pipe(takeUntil(unsubscribe$)).subscribe({
      next: handleMessage
    });
    getWelcomeMessage();

    return () => {
      unsubscribe$.next();
      unsubscribe$.complete();
    };
  }, [subject, getWelcomeMessage, handleMessage]);

  return (
    <div>
      <h3>Batty Chat Bot</h3>
    </div>
  );
}

Home.propTypes = {};

export default Home;

替代方案:用RxJS操作符直接处理流式拼接

既然用了RxJS,可以把消息拼接逻辑放在流里处理,减少React状态的依赖问题:

import { useContext, useEffect, useState } from "react";
import { Context } from "../context/Context";
import { buffer, filter, map } from "rxjs/operators";

function Home() {
  const { subject, getWelcomeMessage } = useContext(Context);
  const [messages, setMessages] = useState([]);

  useEffect(() => {
    // 判断是否为结束标记
    const isEOM = (val) => 
      val.m === "$$$EOM$$$" || val.m === "$$$EOM:CRITICAL$$$" || val.m === "$$$EOM:ERROR$$$";

    const subscription = subject.pipe(
      // 收集直到EOM的所有消息片段
      buffer(subject.pipe(filter(isEOM))),
      // 过滤空片段组
      filter(segments => segments.length > 0),
      // 拼接成完整消息
      map(segments => ({
        role: "assistant",
        content: segments.map(s => s.m).join(''), // 若需要保留数组形式,直接用segments.map(s => s.m)
        inProgress: false
      }))
    ).subscribe(completeMsg => {
      setMessages(prev => [...prev, completeMsg]);
    });

    getWelcomeMessage();
    return () => subscription.unsubscribe();
  }, [subject, getWelcomeMessage]);

  return (
    <div>
      <h3>Batty Chat Bot</h3>
      {messages.map((msg, idx) => (
        <div key={idx}>
          <div>{msg.role}</div>
          <div>{msg.content}</div>
        </div>
      ))}
    </div>
  );
}

Home.propTypes = {};

export default Home;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 10:25:42