Angular+RxJS+DeepChat:解决最后一次发射导致消息重复问题
修复方案
问题根源
当前使用combineLatest操作符会导致重复触发逻辑:该操作符会在任一源Observable发射新值时,将所有源的最新值组合输出。当第二条消息通过messages()发射时,streamPending$的最新值仍是true,会再次生成[true, 第二条消息]的组合结果,导致同一消息被重复处理,最终UI中出现重复消息。
同时原代码的takeWhile(([bool]) => !bool)会在bool为false时直接终止流,无法处理最后一次的完成状态(需要调用signals.onResponse())。
修复代码
使用zip操作符替代combineLatest,确保streamPending$的状态与messages()的消息一一配对,每个消息仅被处理一次;同时调整takeWhile参数以包含最后一次的完成状态,统一逻辑处理:
zip(this.service.streamPending$.pipe(skip(1)), this.messages()).pipe( tap(([isPending, message]) => { // 所有消息都通过addMessage手动添加到DeepChat (this.chatElementRef()?.nativeElement as any).addMessage({ text: message, role: 'ai' }); // 当所有消息接收完成(isPending为false),调用onResponse触发UI完成加载 if (!isPending) { signals.onResponse(); } }), // 第二个参数设为true,确保包含最后一次isPending为false的发射 takeWhile(([isPending]) => isPending, true) ).subscribe();
关键调整说明
zip替代combineLatest:zip会严格按发射顺序将两个Observable的对应位置的值配对,避免因单个源的新发射重复组合旧值,从根源解决消息重复问题。- 统一消息处理逻辑:将所有消息的
addMessage调用整合到一处,避免分散在tap和subscribe中导致逻辑混乱。 - 修正
takeWhile参数:添加第二个参数true,让流在终止时包含isPending为false的那次发射,确保最终能调用signals.onResponse()完成UI状态更新。
内容的提问来源于stack exchange,提问作者DRNR
相关产品推荐
相关产品推荐

