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

如何用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
                )
            }
    }

方案优势

  1. 无时间依赖:不用再纠结生产环境的延迟参数,不会出现同组数据被拆分或处理延迟过高的问题
  2. 性能最大化:本地测试的50万条/秒性能在生产环境可安全落地,无需担心时间窗口的等待损耗
  3. 逻辑精准:基于排序后的连续分组特性,完美捕获完整分组数据,避免漏处理或重复处理

内容的提问来源于stack exchange,提问作者Roar S.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:23:18