Spring AMQP RPC consumer and throw exception

Viewed 1654

I have a consumer (RabbitListner) in RPC mode and I would like to know if it is possible to throw exception that can be treated by the publisher.

To make more clear my explication the case is as follow :

  • The publisher send a message in RPC mode
  • The consumer receive the message, check the validity of the message and if the message can not be take in count, because of missing parameters, then I would like to throw Exception. The exception can be a specific business exception or a particular AmqpException but I want that the publisher can handle this exception if it is not go in timeout.

I try with the AmqpRejectAndDontRequeueException, but my publisher do not receive the exception, but just a response which is empty.

Is it possible to be done or may be it is not a good practice to implement like that ?

EDIT 1 :

After the @GaryRussel response here is the resolution of my question:

  1. For the RabbitListner I create an error handler :

     @Configuration
     public class RabbitErrorHandler implements         RabbitListenerErrorHandler {
     @Override public Object handleError(Message message, org.springframework.messaging.Message<?> message1, ListenerExecutionFailedException e) {
    throw e;
    }
    

    }

  2. Define the bean into a configuration file :

    @Configuration public class RabbitConfig extends RabbitConfiguration {

    @Bean
    public RabbitTemplate getRabbitTemplate() {
        Message.addWhiteListPatterns(RabbitConstants.CLASSES_TO_SEND_OVER_RABBITMQ);
    return new RabbitTemplate(this.connectionFactory());
    

    }

    /**
    * Define the RabbitErrorHandle
    * @return Initialize RabbitErrorHandle bean
    */
    @Bean
    public RabbitErrorHandler rabbitErrorHandler() {
       return new RabbitErrorHandler();
    }
    }
    
  3. Create the @RabbitListner with parameters where rabbitErrorHandler is the bean that I defined previously :

    @Override
    @RabbitListener(queues = "${rabbit.queue}"
        , errorHandler = "rabbitErrorHandler"
        , returnExceptions = "true")
     public ReturnObject receiveMessage(Message message) {
    
  4. For the RabbitTemplate I set this attribute :

       rabbitTemplate.setMessageConverter(new RemoteInvocationAwareMessageConverterAdapter());
    

When the messsage threated by the consumer, but it sent an error, I obtain a RemoteInvocationResult which contains the original exception into e.getCause().getCause().

3 Answers

See the returnExceptions property on @RabbitListener (since 2.0). Docs here.

The returnExceptions attribute, when true will cause exceptions to be returned to the sender. The exception is wrapped in a RemoteInvocationResult object.

On the sender side, there is an available RemoteInvocationAwareMessageConverterAdapter which, if configured into the RabbitTemplate, will re-throw the server-side exception, wrapped in an AmqpRemoteException. The stack trace of the server exception will be synthesized by merging the server and client stack traces.

Important

This mechanism will generally only work with the default SimpleMessageConverter, which uses Java serialization; exceptions are generally not "Jackson-friendly" so can’t be serialized to JSON. If you are using JSON, consider using an errorHandler to return some other Jackson-friendly Error object when an exception is thrown.

What worked for me was :

On "serving" side :

  • Service

    @RabbitListener(id = "test1", containerFactory ="BEAN CONTAINER FACTORY", 
    queues = "TEST QUEUE", returnExceptions = "true")
    DataList getData() {
      // this exception will be transformed by rabbit error handler to a RemoteInvocationResult
     throw new IllegalStateException("mon expecion");
     //return dataHelper.loadAllData();
    }
    

On "requesting" side :

  • Service

    public void fetchData() throws AmqpRemoteException {
      var response = (DataList) amqpTemplate.convertSendAndReceive("TEST EXCHANGE", "ROUTING NAME", new Object());
    
      Optional.ofNullable(response)
             .ifPresentOrElse(this::setDataContent, this::handleNoData);
    }
    
  • Config

    @Bean
    AmqpTemplate amqpTemplate(ConnectionFactory connectionFactory, MessageConverter messageConverter) {
       var rabbitTemplate = new RabbitTemplate(connectionFactory);
       rabbitTemplate.setMessageConverter(messageConverter);
       return rabbitTemplate;
    }
    @Bean
    MessageConverter jsonMessageConverter() {
       ObjectMapper objectMapper = new ObjectMapper();
       objectMapper.configure(SerializationFeature.FAIL_ON_EMPTY_BEANS, false);
       objectMapper.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
       objectMapper.registerModule(new JavaTimeModule());
    
       var jsonConverter = new Jackson2JsonMessageConverter(objectMapper);
    
       DefaultClassMapper classMapper = new DefaultClassMapper();
       Map<String, Class<?>> idClassMapping = Map.of(
          DataList.class.getName(), DataList.class,
          RemoteInvocationResult.class.getName(), RemoteInvocationResult.class
       );
       classMapper.setIdClassMapping(idClassMapping);
    
       jsonConverter.setClassMapper(classMapper);
    
       // json converter with returned exception awareness
       // this will transform RemoteInvocationResult into a AmqpRemoteException 
       return new RemoteInvocationAwareMessageConverterAdapter(jsonConverter);
    }
    

You have to return a message as an error, which the consuming application can choose to treat as an exception. However, I don't think normal exception handling flows apply with messaging. Your publishing application (the consumer of the RPC service) needs to know what can go wrong and be programmed to deal with those possibilities.

Related