Why can't i subscribe multiple times to a publisher?

Viewed 130

i wrote two methods,first one calls the publisher (productDtoMono) only once and the other one calls it twice, the first one returns bad request (empty body) when called from controller while the second works perfectly and returns http ok i did my homework and i see that i can't consume the publisher more than one time why is that ?

public Mono<ProductDto> firstupdateProduct(Mono<ProductDto> productDtoMono) {
    return productDtoMono.flatMap(productDto ->
            productRepository.findById(productDto.getId())
                    .flatMap(p -> productDtoMono
                            .map(EntityDtoUtil::toEntity)))
                    .flatMap(productRepository::save)
                    .map(EntityDtoUtil::toDto);
}
public Mono<ProductDto> secondupdateProduct(String id, Mono<ProductDto> productDtoMono) {
    return productRepository.findById(id)
            .flatMap(product -> productDtoMono
                    .map(EntityDtoUtil::toEntity))
            .flatMap(productRepository::save)
            .map(EntityDtoUtil::toDto);
}

my controller method :

@PutMapping
public Mono<ResponseEntity<ProductDto>> updateProduct(@RequestBody Mono<ProductDto> productDtoMono){
    System.out.println("in controller "+productDtoMono);
    return productService.updateProduct(productDtoMono)
            .map(ResponseEntity::ok)
            .defaultIfEmpty(ResponseEntity.notFound().build());

}
1 Answers

I did more research. Mono, Flux are cold publishers by default, so when the source of their creation comes from a request, they can only be consumed once. When we try to consume it a second time, it tries to generate the object from a request, hence the bad request get thrown. I had to cache the body while passing it to service like this:

@PutMapping
public Mono<ResponseEntity<ProductDto>> updateProduct(@RequestBody Mono<ProductDto> productDtoMono){
    System.out.println("in controller "+productDtoMono);
    return productService.updateProduct(productDtoMono.cache())
            .map(ResponseEntity::ok)
            .defaultIfEmpty(ResponseEntity.notFound().build());

}

and changed updateProduct like this

public Mono<ProductDto> updateProduct(Mono<ProductDto> productDtoMono) {
    return productDtoMono.map(ProductDto::getId)
            .flatMap(productRepository::findById)
            .flatMap(product -> productDtoMono.map(EntityDtoUtil::toEntity))
            .flatMap(productRepository::save)
            .map(EntityDtoUtil::toDto);
}
Related