Flux: How to prevent element processing when some element has concrete value?

Viewed 277

The following test method:

@Test
void testMe() {
    Flux.just(1, 2, 3, 4, 5)
            .map(this::saveInDb)
            .toStream().count();
}

int saveInDb(int element) {
    System.out.println(element + " successfully stored in DB.");
    return element;
}

prints always:

1 successfully stored in DB.
2 successfully stored in DB.
3 successfully stored in DB.
4 successfully stored in DB.
5 successfully stored in DB.

The question that I have: how to prevent any saving to DB when some particular element available in the Flux?

For example: I do not want to store anything to DB when element 5 exists in Flux.

Is it even possible to implement this requirement in a non-blocking fashion?

1 Answers

It is only possible if you can fit the whole Flux into memory:

Flux.just(1, 2, 3, 4, 5)
    .collectList()
    .filter(items -> !items.contains(5))
    .flatMapIterable(x -> x)
    .map(this::saveInDb)

If the Flux can't fit into the memory you can:

  • buffer the Flux and check the presence only in the buffered subsets (this might not be right for your use case)
  • consume the Flux twice, first check the presence with a filter, if the item was not present then subscribe again (this obviously only works if the Flux is guaranteed to have the same output on both occasions)
Related