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

RxJS如何组合两个Observable实现离线缓冲联网触发数据同步

RxJS 离线/在线全量数据同步实现

需求对应逻辑

需要组合两个基础Observable实现同步规则:

  • connection$:网络状态流,状态变更时推送布尔值,true为在线、false为离线
  • data$:待同步数据流,产生新待同步数据时推送对应值

组合规则要求:

  • 离线状态下不推送任何值,全量缓存data$产生的所有待同步数据,不能只保留最新值
  • 在线状态下data$产生新值时立即推送
  • 离线切回在线的瞬间,按数据产生顺序推送所有离线阶段缓存的数据,之后回到在线实时推送模式

之前用filter+combineLatest实现不符合预期的核心原因是:combineLatest只会持有每个输入流的最新值,本身没有全量缓冲能力,也无法实现离线阶段的推送拦截,完全不匹配当前场景。

可直接使用的实现代码

核心通过连接状态做缓冲开关,拆分离线缓存释放、在线实时推送两个通道合并实现,不需要轮询,完全基于RxJS标准组合逻辑:

import { merge } from 'rxjs';
import { buffer, filter, map, mergeMap, pairwise, startWith, withLatestFrom } from 'rxjs/operators';

/**
 * 生成离线在线同步流
 * @param {Observable<boolean>} connection$ 连接状态流
 * @param {Observable<any>} data$ 待同步数据流
 * @returns {Observable<any>} 同步触发流
 */
function createSyncStream(connection$, data$) {
  // 离线切在线的事件流:仅在状态从false变为true时触发
  const onlineSwitch$ = connection$.pipe(
    startWith(false), // 初始默认离线,避免订阅时已在线的场景误触发缓存
    pairwise(),
    filter(([prevStatus, currStatus]) => prevStatus === false && currStatus === true)
  );

  // 离线缓存通道:缓冲所有离线期间产生的数据,上线时依次释放所有缓存值
  const offlineBuffered$ = data$.pipe(
    buffer(onlineSwitch$),
    mergeMap(cachedDataList => cachedDataList) // 把缓存数组拆成单个值按序推送
  );

  // 在线实时通道:仅当当前状态为在线时,直接透传新产生的数据
  const onlineRealtime$ = data$.pipe(
    withLatestFrom(connection$),
    filter(([_, isOnline]) => isOnline),
    map(([data]) => data)
  );

  // 合并两个通道即为最终同步流
  return merge(offlineBuffered$, onlineRealtime$);
}

实现说明

  • 无数据丢失:离线阶段所有data$推送的值都会被buffer完整收集,直到上线瞬间一次性按产生顺序输出
  • 无重复推送:离线缓存通道仅在上线切换瞬间触发一次缓存释放,在线阶段新数据全部走实时通道,两个通道不会产生重复值
  • 可扩展:如果需要限制缓存大小、添加同步失败重试逻辑、数据去重,直接在对应通道的pipe中添加操作符即可,不需要改动核心逻辑。

内容的提问来源于stack exchange,提问作者Gernot R. Bauer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 10:09:15