如何从Flux中按name分组选取优先级最高的元素,去除同属性重复对象
实现方案
你使用distinct(MyObject::getName)无法满足需求是因为该操作符默认保留对应key首次出现的元素,不支持自定义保留规则。针对按name去重且保留priority最大值的需求,有限流场景下的最优实现如下:
核心处理代码
Flux<MyObject> resultFlux = myFlux // 按name属性分组,同name的元素会被分到同一个组 .groupBy(MyObject::getName) // 每个分组内仅保留priority最大的元素 .flatMap(group -> group.reduce((prev, curr) -> prev.getPriority() > curr.getPriority() ? prev : curr ));
完整可运行示例
import lombok.Data; import reactor.core.publisher.Flux; public class Example { @Data // 调整为静态内部类方便直接实例化,不影响核心逻辑 public static class MyObject { final String name; final int priority; } public static void main(String[] args) { Flux<MyObject> myFlux = Flux.just( new MyObject("abc", 2), new MyObject("abc", 4), new MyObject("cde", 1)); Flux<MyObject> resultFlux = myFlux .groupBy(MyObject::getName) .flatMap(group -> group.reduce((prev, curr) -> prev.getPriority() > curr.getPriority() ? prev : curr )); // 验证输出,结果和预期完全一致 resultFlux.collectList().subscribe(System.out::println); // 输出:[Example.MyObject(name=abc, priority=4), Example.MyObject(name=cde, priority=1)] } }
方案说明
- 完全基于Reactor原生操作符实现,无额外依赖,符合响应式编程规范
- 有限流场景下时间复杂度为O(n),仅需遍历一次流元素,性能最优
- 支持任意顺序的流元素,不需要提前按name排序
如果是无限流场景,需要搭配窗口操作符按时间/元素数量触发去重,避免分组内存泄漏,示例如下:
// 每处理100个元素触发一次当前窗口内的去重 Flux<MyObject> infiniteResultFlux = myFlux .window(100) .flatMap(window -> window .groupBy(MyObject::getName) .flatMap(group -> group.reduce((prev, curr) -> prev.getPriority() > curr.getPriority() ? prev : curr )) );
内容的提问来源于stack exchange,提问作者DaviF
相关产品推荐
相关产品推荐

