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

Flowable流拆分处理差异后重组时更新异常问题求助

--- Stream A [ A.map1.diff() ] ---
          |                                     |
          |                                     |
Source --- --- Stream B [ B.map5.diff() ] --- --- combineLatest(A,B,C)
          |                                     |
          |                                     |
           --- Stream C [ C.map9.diff() ] ---

Source始终会发送所有map的当前完整值,各个子流仅需发送对应map的差异值。diff是已在其他场景验证可用的Flowable扩展方法,无差异时不会发送数据。

当map9出现差异时,combineLatest会包含map9的更新值,但map1和map5仍保持流创建时的初始完整状态(因为需要先获取完整值才能计算差异);当map5出现差异时,map1和map9也会出现同样的情况。

这就导致combineLatest每次发送的更新数据量过大——本质上除了触发更新的map,其他都是原始完整值。我尝试过相关的Rx流拆分与组合方案,但未能解决问题。


补充说明

  • 这是不可修改的Socket连接结构的一部分:存在所有连接共享的根Flowable,以及每个连接对应的独立订阅Flowable。
  • 不一定要拆分Flowable流,如果有无需拆分、直接逐个修改FlowAgg字段的方法,也可以采用。

代码片段

data class FlowAgg(
    val devices: Map<Int, Device>,
    val assignments: Map<Int, Assignment>,
    val systemtime: Map<Int, Timestamp>
)
data class Summary(
    val id: Int,
    val device: Device? = null,
    val assignment: Assignment? = null,
    val systemtime: Timestamp? = null
)

[...]

socketTopic(
    path = "/summary",
    root = { _ ->
        Flowables.combineLatest(
            DeviceFlowable,
            AssignmentFlowable,
            SystemtimeFlowable
        ) { devices, assignments, systemtime ->
            FlowAgg(
                devices = devices,
                assignments = assignments,
                systemtime = systemtime,
            )
        }
    },
    subscription = { broadcast ->
        broadcast
            .publish { flow ->
                Flowables.combineLatest(
                    flow.map { it.devices }.diff(),
                    flow.map { it.assignments }.diff(),
                    flow.map { it.systemtime }.diff()
                ) { devicesDiff, assignmentsDiff, systemTimeDiff ->
                    val keys = devicesDiff.keys + assignmentsDiff.keys + systemTimeDiff.keys
                    keys.map { id ->
                        Summary(
                            id = id,
                            device = devicesDiff[id],
                            assignment = assignmentsDiff[id],
                            systemtime = systemTimeDiff[id]
                        )
                    }
                }
                .map {
                    Json.encodeToString(ListSerializer(Summary.serializer()), it)
                }
            }
    }
)

(注:修正了原代码中Summary字段拼写错误assigment为assignment,以及赋值时的错误引用)


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 00:01:20