【问题标题】:Auto Delete Messages from Queue Once consumed using Spring AMQP使用 Spring AMQP 使用后自动从队列中删除消息
【发布时间】:2016-08-27 06:00:38
【问题描述】:

我有 2 个应用程序使用 RabbitMQ 交换数据。我已经使用 Spring AMQP 实现了这一点。我有一个场景,一旦从消费者那里消费了消息,可能会在处理时遇到异常。

如果出现任何异常,我计划登录数据库。一旦消息到达消费者,无论是成功处理还是遇到错误,我都必须明确地从队列中删除消息。

如何从队列中强制删除消息,否则它将是 如果我的申请无法处理呢?

下面是我的监听器代码

 @RabbitListener(containerFactory="rabbitListenerContainerFactory",queues=Constants.JOB_QUEUE)
            public void handleMessage(JobListenerDTO jobListenerDTO) {
                //System.out.println("Received summary: " + jobListenerDTO.getProcessXML());
                //amqpAdmin.purgeQueue(Constants.JOB_QUEUE, true);
                try{
                    Map<String, Object> variables = new HashMap<String, Object>();  
                    variables.put("initiator", "cmy5kor");

                    Deployment deploy = repositoryService.createDeployment().addString(jobListenerDTO.getProcessId()+".bpmn20.xml",jobListenerDTO.getProcessXML()).deploy();
                    ProcessInstance processInstance = runtimeService.startProcessInstanceByKey(jobListenerDTO.getProcessId(), variables);

                    System.out.println("Process Instance is:::::::::::::"+processInstance);

                }catch(Exception e){

                    e.printStackTrace();
            }

配置代码

@Configuration
@EnableRabbit
public class RabbitMQJobConfiguration extends AbstractBipRabbitConfiguration {


    @Bean
    public RabbitTemplate rabbitTemplate() {
        RabbitTemplate template = new RabbitTemplate(connectionFactory());
        template.setQueue(Constants.JOB_QUEUE);
        template.setMessageConverter(jsonMessageConverter());
        return template;
    }


    @Bean
    public Queue jobQueue() {
        return new Queue(Constants.JOB_QUEUE);
    }


    @Bean(name="rabbitListenerContainerFactory")
    public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory() {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory());
        Jackson2JsonMessageConverter messageConverter = new Jackson2JsonMessageConverter();
        DefaultClassMapper classMapper = new DefaultClassMapper();
        Map<String, Class<?>> idClassMapping = new HashMap<String, Class<?>>();
        idClassMapping.put("com.bosch.diff.approach.TaskMessage", JobListenerDTO.class);
        classMapper.setIdClassMapping(idClassMapping);
        messageConverter.setClassMapper(classMapper);
        factory.setMessageConverter(messageConverter);
        factory.setReceiveTimeout(10L);
        return factory;
    }



}

【问题讨论】:

    标签: java spring rabbitmq spring-amqp


    【解决方案1】:

    我不知道 spring api 或 rmq 的配置,但是这个

    一旦消息到达消费者,无论是成功处理还是遇到错误,我都必须从队列中显式删除消息。

    正是您设置自动确认标志时发生的情况。通过这种方式,消息在被消费后立即得到确认 - 因此从队列中消失。

    【讨论】:

      【解决方案2】:

      只要您的侦听器捕获到异常,消息就会从队列中移除。

      如果你的监听器抛出异常,默认会重新入队;可以通过抛出 AmqpRejectAndDontRequeueException 或设置 defaultRequeueRejected 属性来修改该行为 - 请参阅 the documentation for details。

      【讨论】:

      • 谢谢加里。 我们在 github 中是否有一些示例供我阅读和实施?
      • 见the samples。
      • @GaryRussell,据我所知,rabbitmq 在得到 ACK 后从队列中删除消息,是吗?那么当 spring-boot 发送 ACK 时?如果任意(AmqpRejectAndDontRequeueException 除外)未捕获的异常来自侦听器?好的,但是如果我自己的转换器出现未处理的异常(例如解析器异常)怎么办?
      • 查看我对your question的回复。
      • 您好,能否更新一下链接。它说无效。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-21
      • 1970-01-01
      • 2016-12-13
      • 2018-08-04
      • 1970-01-01
      • 2011-03-26
      相关资源
      最近更新 更多