How to select max item per group from Flux

Viewed 891

Given the following MyObject and Flux<MyObject> what is the best way to remove MyObjects with same the same property from this flux?

    import lombok.Data;
    import reactor.core.publisher.Flux;
    
    public class Example {
    
        @Data
        public class MyObject {
            final String name;
            final int priority;
        }
    
        public Example() {
            Flux<MyObject> myFlux = Flux.just(
                    new MyObject("abc", 2),
                    new MyObject("abc", 4),
                    new MyObject("cde", 1));
        }
    }

For example I want to remove objects with same name while choosing the ones with higher priority.
Output: [Example.MyObject(name=abc, priority=4), Example.MyObject(name=cde, priority=1)]

If I use myFlux.distinct(MyObject::getName) I wont be able to choose which one to keep.

2 Answers

You can achieve this using groupBy and reduce operators on Flux:

Flux.just(
        new MyObject("abc", 2),
        new MyObject("abc", 4),
        new MyObject("cde", 1))
    .groupBy(MyObject::getName)
    .flatMap(group -> group.reduce((o1, o2) -> o1.getPriority() > o2.getPriority() ? o1 : o2))
    .subscribe(System.out::println);

One important consideration is that this only works well if the number of groups is small, otherwise it can result in a deadlock. As a remedy you can set the maxConcurrency parameter of flatMap to a higher value.

See documentation of groupBy operator:

The groups need to be drained and consumed downstream for groupBy to work correctly. Notably when the criteria produces a large amount of groups, it can lead to hanging if the groups are not suitably consumed downstream (eg. due to a flatMap with a maxConcurrency parameter that is set too low).

To solve this, you first need to convert this Flux<MyObject> into a Mono<List<MyObject>> because you need to know all objects and their priorities in order to sort them.

As soon as you have the list of all instances of MyObject, you can use the Java 8 Stream api to solve this:

@Slf4j
public class Example {

    public static void main(String[] args) {
        Flux<MyObject> myFlux = Flux.just(
                        new MyObject("abc", 2),
                        new MyObject("abc", 4),
                        new MyObject("cde", 1))
                .collectList()
                .map(myObjectsList -> myObjectsList.stream()
                        .collect(Collectors
                                .groupingBy(MyObject::getName)))
                // now we have a Map<String, List<MyObject>>
                .map(Map::entrySet)
                // now we have a Set<Entry<String, List<MyObject>>>
                .flatMapIterable(entrySet -> entrySet)
                .map(Map.Entry::getValue)
                // now we have a Flux<List<MyObject>>
                // and all MyObject in that list have 
                // the same name
                .filter(allObjectsWithSameName -> !allObjectsWithSameName.isEmpty())
                // now we sort all the lists in descending order
                // and return the first element
                // which is the one with the highest prio
                .map(allObjectsWithSameName -> {
                            allObjectsWithSameName.sort(new Comparator<MyObject>() {
                                @Override
                                public int compare(MyObject o1, MyObject o2) {
                                    return Integer.compare(o2.priority, o1.priority);
                                }
                            });
                            return allObjectsWithSameName.get(0);
                        }
                );

        myFlux.subscribe(result -> System.out.println("MyObject: " + result.toString()));
    }

    @Data
    @RequiredArgsConstructor
    public static class MyObject {
        final String name;
        final int priority;
    }
}

Output:

MyObject: Example.MyObject(name=abc, priority=4)
MyObject: Example.MyObject(name=cde, priority=1)
Related