【问题标题】:Spring Integration JMS assured message delivery using DSLSpring Integration JMS 使用 DSL 确保消息传递
【发布时间】:2018-06-26 13:31:20
【问题描述】:

我正在尝试创建一个流(1),其中从 TCP 适配器接收消息,该适配器可以是客户端或服务器,并将消息发送到 ActiveMQ 代理。

我的另一个流程(2)从所需队列中挑选消息并发送到目的地

TCP(客户端/服务器)==(1)==> ActiveMQ 代理 ==(2)==> HTTP 出站适配器

我想确保如果我的消息没有传递到所需的目的地,那么它会重新尝试再次发送消息。

我当前流向代理的流程 (1) 是:

IntegrationFlow flow = IntegrationFlows
            .from(Tcp
                    .inboundAdapter(Tcp.netServer(Integer.parseInt(1234))
                            .serializer(customSerializer).deserializer(customSerializer)
                            .id("server").soTimeout(5000))
                    .id(hostConnection.getConnectionNumber() + "adapter"))).channel(directChannel())
            .wireTap("tcpInboundMessageLogChannel").channel(directChannel())
            .handle(Jms.outboundAdapter(activeMQConnectionFactory)
                    .destination("jmsInbound"))
            .get();

    this.flowContext.registration(flow).id("outflow").register();

和我从代理到 http 出站的流程(2):

flow = IntegrationFlows
            .from(Jms.messageDrivenChannelAdapter(activeMQConnectionFactory)
                    .destination("jmsInbound"))
            .channel(directChannel())
            .handle(Http.outboundChannelAdapter(hostConnection.getUrl()).httpMethod(HttpMethod.POST)
                    .expectedResponseType(String.class)
                    .mappedRequestHeaders("abc"))
            .get();
    this.flowContext.registration(flow).id("inflow").register();

问题:

  • 如果在传递过程中出现任何异常,例如我的目标 URL 不起作用,那么它会重新尝试发送消息。

  • 尝试失败后重试7次,即max attempt to 7

  • 如果尝试仍然不成功,则将消息发送到ActiveMQ.DLQ(死信队列)并且不会再次尝试,因为消息从实际队列中出列并发送到ActiveMQ.DLQ。

所以,我希望没有消息丢失并且消息将按顺序处理的场景。

【问题讨论】:

    标签: spring-integration spring-jms spring-integration-dsl


    【解决方案1】:

    首先:我相信你可以为无限重试配置jmsInbound:

    /**
     * Configuration options for a messageConsumer used to control how messages are re-delivered when they
     * are rolled back.
     * May be used server side on a per destination basis via the Broker RedeliveryPlugin
     *
     * @org.apache.xbean.XBean element="redeliveryPolicy"
     *
     */
    public class RedeliveryPolicy extends DestinationMapEntry implements Cloneable, Serializable {
    

    另一方面,您可以为 RequestHandlerRetryAdvice 配置一个 .handle(Http.outboundChannelAdapter( 以实现类似的重试行为,但在应用程序内部无需往返 JMS 并返回:https://docs.spring.io/spring-integration/docs/5.0.6.RELEASE/reference/html/messaging-endpoints-chapter.html#retry-advice

    以下是一些如何从 Java DSL 角度对其进行配置的示例:

        @Bean
        public IntegrationFlow errorRecovererFlow() {
            return IntegrationFlows.from(Function.class, "errorRecovererFunction")
                    .handle((GenericHandler<?>) (p, h) -> {
                        throw new RuntimeException("intentional");
                    }, e -> e.advice(retryAdvice()))
                    .get();
        }
    
        @Bean
        public RequestHandlerRetryAdvice retryAdvice() {
            RequestHandlerRetryAdvice requestHandlerRetryAdvice = new RequestHandlerRetryAdvice();
            requestHandlerRetryAdvice.setRecoveryCallback(new ErrorMessageSendingRecoverer(recoveryChannel()));
            return requestHandlerRetryAdvice;
        }
    
        @Bean
        public MessageChannel recoveryChannel() {
            return new DirectChannel();
        }
    

    RequestHandlerRetryAdvice 可以与RetryTemplate 一起配置以应用类似AlwaysRetryPolicy 的内容。有关更多信息,请参阅 Spring Retry 项目:https://github.com/spring-projects/spring-retry

    【讨论】:

      猜你喜欢
      • 2012-06-30
      • 2016-02-10
      • 1970-01-01
      • 2019-06-29
      • 1970-01-01
      • 2020-04-30
      • 2014-09-27
      • 2017-12-12
      • 2016-05-07
      相关资源
      最近更新 更多