如何用takeUntil替代Flux.groupBy中的固定延迟take方法?
问题描述
我用Flux的groupBy方法结合.take(Duration.ofMillis(10))处理数据,每秒可处理约5万条记录。本地用Flux.just测试时,将延迟设为1ms能达到每秒50万条的处理量,但不敢在生产环境中这么做。
数据来自BigQuery数据库,且已按分组字段排序。请问是否可以避免使用带固定延迟的.take(Duration.ofMillis(10)),转而使用takeUntil这类方法?
原代码
data class IncomingDto( val aggregation: Aggregation, val municipalityId: String, val personId: String, val startDate: LocalDate, val endDate: LocalDate? = null, ) data class ResultDto( val aggregation: Aggregation, val municipalityId: String, val itemCount: Int, ) fun reduceResult( someDtoFlux: Flux<IncomingDto> ): Flux<ResultDto> = someDtoFlux .groupBy { it.aggregation } // 生产代码中按4个字段分组 .flatMap { groupFlux -> groupFlux .take(Duration.ofMillis(10)) // 依赖固定时间窗口 .collectList() .mapNotNull { currentSomeDtoList -> // 此处为归约逻辑 // 模拟结果 ResultDto( aggregation = currentSomeDtoList.first().aggregation, municipalityId = currentSomeDtoList.first().municipalityId, itemCount = 42 ) } }
解决方案
当然可以,而且这才是适配你场景的最优方案——你的数据已经按分组字段排序,同组记录必然连续出现,完全可以利用这个特性替代固定延迟的时间窗口,彻底摆脱时间依赖。
核心思路
用bufferUntilChanged或windowUntilChanged替代groupBy+固定延迟take的组合:
bufferUntilChanged会把连续相同分组的记录打包成一个List,直到分组字段变化时输出- 无需等待时间窗口,数据凑齐就处理,性能和准确性都能拉满
改造后代码(基于bufferUntilChanged)
data class IncomingDto( val aggregation: Aggregation, val municipalityId: String, val personId: String, val startDate: LocalDate, val endDate: LocalDate? = null, ) data class ResultDto( val aggregation: Aggregation, val municipalityId: String, val itemCount: Int, ) fun reduceResult( someDtoFlux: Flux<IncomingDto> ): Flux<ResultDto> = someDtoFlux // 生产环境替换为4个分组字段的组合(比如用Pair或自定义数据类) .bufferUntilChanged { dto -> dto.aggregation } .mapNotNull { currentGroupList -> // 复用你原有的归约逻辑即可 ResultDto( aggregation = currentGroupList.first().aggregation, municipalityId = currentGroupList.first().municipalityId, itemCount = currentGroupList.size // 示例用实际数量,可替换为你的业务逻辑 ) }
内存优化版(基于windowUntilChanged)
如果处理的分组数据量极大,用windowUntilChanged可以避免一次性把整个分组加载到内存,而是流式处理:
fun reduceResult( someDtoFlux: Flux<IncomingDto> ): Flux<ResultDto> = someDtoFlux .windowUntilChanged { it.aggregation } .flatMap { windowFlux -> // 流式统计分组内的数量(替换为你的归约逻辑) val countFlux = windowFlux.reduce(0) { acc, _ -> acc + 1 } // 获取分组的标识信息(取第一个元素即可) val firstDtoFlux = windowFlux.take(1) countFlux.zipWith(firstDtoFlux) .map { (itemCount, firstDto) -> ResultDto( aggregation = firstDto.aggregation, municipalityId = firstDto.municipalityId, itemCount = itemCount ) } }
方案优势
- 无时间依赖:不用再纠结生产环境的延迟参数,不会出现同组数据被拆分或处理延迟过高的问题
- 性能最大化:本地测试的50万条/秒性能在生产环境可安全落地,无需担心时间窗口的等待损耗
- 逻辑精准:基于排序后的连续分组特性,完美捕获完整分组数据,避免漏处理或重复处理
内容的提问来源于stack exchange,提问作者Roar S.
相关产品推荐
相关产品推荐

