[ 
https://issues.apache.org/jira/browse/CAMEL-12906?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

michael elbaz updated CAMEL-12906:
----------------------------------
    Description: 
/!\ Its about camel aggregator and error handling during *shutdown phase*

My case is about graceful shutdown and error handling in case of aggregator 
first one:

In this example i will get *RejectedExecutionException* during the shutdown and 
the exception is not catched by the *onexception* or *errorHandler* methods as 
expected so i will just lost message.

{code:java}
@Component
public class AmqRoute extends RouteBuilder {

    @Override
    public void configure() throws Exception {
        errors();

        from("timer:foo?period=1000")
                .log("1")
                .transform().body(() -> "DATA " + 
RandomStringUtils.randomAlphanumeric(10))
                .convertBodyTo(String.class)
                .to(amq());

        from(amq()).to("direct:foo");

        from("direct:foo")
                .delay(3000)
                .aggregate(constant(true), new GroupedBodyAggregationStrategy())
                .completionSize(10)
                .forceCompletionOnStop()
                .log("${body}");
    }

    public void errors() {
        onException(Exception.class)
                .useOriginalMessage()
                .to("activemq:recovery")
                .handled(true)
                .onRedelivery(exchange -> System.err.println("push to amq"));

        errorHandler(deadLetterChannel("log:dead?level=ERROR"));
    }

    private static String amq() {
        String amq = "activemq:data";

        amq += "?transacted=true";

        return amq;
    }

}
{code}

In this one i'll will catch the exception using doTry and the doCatch part will 
never be reached (i read the doc about but it just weird the fact we can put 
some unreachable doCatch and that the only way to make a doTry on aggregate and 
it's look there is no exception throwed, that the interesting point because i 
expect at least catched *RejectedExecutionException* and finally i *didn't 
lost* message in this case !

{code:java}
@Component
public class AmqRoute extends RouteBuilder {

    @Override
    public void configure() throws Exception {
       // errors();

        from("timer:foo?period=1000")
                .log("1")
                .transform().body(() -> "DATA " + 
RandomStringUtils.randomAlphanumeric(10))
                .convertBodyTo(String.class)
                .to(amq());

        from(amq()).to("direct:foo");

        from("direct:foo")
                .delay(3000)
                .doTry()
                .aggregate(constant(true), new GroupedBodyAggregationStrategy())
                .completionSize(10)
                .forceCompletionOnStop()
                .log("${body}")
                .endDoTry()
                .doCatch(Exception.class)
                .to("activemq:recovery")
                .log(LoggingLevel.ERROR, "!!!!!!!!!!!!!!!!")
                .end();
    }

    public void errors() {
        onException(Exception.class)
                .useOriginalMessage()
                .to("activemq:recovery")
                .handled(true)
                .onRedelivery(exchange -> System.err.println("push to amq"));

        errorHandler(deadLetterChannel("log:dead?level=ERROR"));
    }

    private static String amq() {
        String amq = "activemq:data";

        amq += "?transacted=true";

        return amq;
    }

}
{code}


  was:
/!\ Its about camel aggregator and error handling during *shutdown phase*

My case is about graceful shutdown and error handling in case of aggregator 
first one:

In this example i will get *RejectedExecutionException* when during the 
shutdown and the exception is not catched by the *onexception* method so i will 
just lost message
{code:java}
@Component
public class AmqRoute extends RouteBuilder {

    @Override
    public void configure() throws Exception {
        errors();

        from("timer:foo?period=1000")
                .log("1")
                .transform().body(() -> "DATA " + 
RandomStringUtils.randomAlphanumeric(10))
                .convertBodyTo(String.class)
                .to(amq());

        from(amq()).to("direct:foo");

        from("direct:foo")
                .delay(3000)
                .aggregate(constant(true), new GroupedBodyAggregationStrategy())
                .completionSize(10)
                .forceCompletionOnStop()
                .log("${body}");
    }

    public void errors() {
        onException(Exception.class)
                .useOriginalMessage()
                .to("activemq:recovery")
                .handled(true)
                .onRedelivery(exchange -> System.err.println("push to amq"));

        errorHandler(deadLetterChannel("log:dead?level=ERROR"));
    }

    private static String amq() {
        String amq = "activemq:data";

        amq += "?transacted=true";

        return amq;
    }

}
{code}

In the seconde one example i'll will catch the exception and the doCatch part 
will never be reached (i read the doc about but it just weird the fact we can 
put some unreachable doCatch and that the only way to make a doTry on aggregate 
and it's seem's the is no exception throwed that the interesting point because 
i excpect at least catched *RejectedExecutionException* and finally i *didn't 
lost* message in this case !

{code:java}
@Component
public class AmqRoute extends RouteBuilder {

    @Override
    public void configure() throws Exception {
       // errors();

        from("timer:foo?period=1000")
                .log("1")
                .transform().body(() -> "DATA " + 
RandomStringUtils.randomAlphanumeric(10))
                .convertBodyTo(String.class)
                .to(amq());

        from(amq()).to("direct:foo");

        from("direct:foo")
                .delay(3000)
                .doTry()
                .aggregate(constant(true), new GroupedBodyAggregationStrategy())
                .completionSize(10)
                .forceCompletionOnStop()
                .log("${body}")
                .endDoTry()
                .doCatch(Exception.class)
                .to("activemq:recovery")
                .log(LoggingLevel.ERROR, "!!!!!!!!!!!!!!!!")
                .end();
    }

    public void errors() {
        onException(Exception.class)
                .useOriginalMessage()
                .to("activemq:recovery")
                .handled(true)
                .onRedelivery(exchange -> System.err.println("push to amq"));

        errorHandler(deadLetterChannel("log:dead?level=ERROR"));
    }

    private static String amq() {
        String amq = "activemq:data";

        amq += "?transacted=true";

        return amq;
    }

}
{code}



> Strange comportement with aggregator
> ------------------------------------
>
>                 Key: CAMEL-12906
>                 URL: https://issues.apache.org/jira/browse/CAMEL-12906
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-activemq
>    Affects Versions: 2.22.1
>            Reporter: michael elbaz
>            Priority: Major
>
> /!\ Its about camel aggregator and error handling during *shutdown phase*
> My case is about graceful shutdown and error handling in case of aggregator 
> first one:
> In this example i will get *RejectedExecutionException* during the shutdown 
> and the exception is not catched by the *onexception* or *errorHandler* 
> methods as expected so i will just lost message.
> {code:java}
> @Component
> public class AmqRoute extends RouteBuilder {
>     @Override
>     public void configure() throws Exception {
>         errors();
>         from("timer:foo?period=1000")
>                 .log("1")
>                 .transform().body(() -> "DATA " + 
> RandomStringUtils.randomAlphanumeric(10))
>                 .convertBodyTo(String.class)
>                 .to(amq());
>         from(amq()).to("direct:foo");
>         from("direct:foo")
>                 .delay(3000)
>                 .aggregate(constant(true), new 
> GroupedBodyAggregationStrategy())
>                 .completionSize(10)
>                 .forceCompletionOnStop()
>                 .log("${body}");
>     }
>     public void errors() {
>         onException(Exception.class)
>                 .useOriginalMessage()
>                 .to("activemq:recovery")
>                 .handled(true)
>                 .onRedelivery(exchange -> System.err.println("push to amq"));
>         errorHandler(deadLetterChannel("log:dead?level=ERROR"));
>     }
>     private static String amq() {
>         String amq = "activemq:data";
>         amq += "?transacted=true";
>         return amq;
>     }
> }
> {code}
> In this one i'll will catch the exception using doTry and the doCatch part 
> will never be reached (i read the doc about but it just weird the fact we can 
> put some unreachable doCatch and that the only way to make a doTry on 
> aggregate and it's look there is no exception throwed, that the interesting 
> point because i expect at least catched *RejectedExecutionException* and 
> finally i *didn't lost* message in this case !
> {code:java}
> @Component
> public class AmqRoute extends RouteBuilder {
>     @Override
>     public void configure() throws Exception {
>        // errors();
>         from("timer:foo?period=1000")
>                 .log("1")
>                 .transform().body(() -> "DATA " + 
> RandomStringUtils.randomAlphanumeric(10))
>                 .convertBodyTo(String.class)
>                 .to(amq());
>         from(amq()).to("direct:foo");
>         from("direct:foo")
>                 .delay(3000)
>                 .doTry()
>                 .aggregate(constant(true), new 
> GroupedBodyAggregationStrategy())
>                 .completionSize(10)
>                 .forceCompletionOnStop()
>                 .log("${body}")
>                 .endDoTry()
>                 .doCatch(Exception.class)
>                 .to("activemq:recovery")
>                 .log(LoggingLevel.ERROR, "!!!!!!!!!!!!!!!!")
>                 .end();
>     }
>     public void errors() {
>         onException(Exception.class)
>                 .useOriginalMessage()
>                 .to("activemq:recovery")
>                 .handled(true)
>                 .onRedelivery(exchange -> System.err.println("push to amq"));
>         errorHandler(deadLetterChannel("log:dead?level=ERROR"));
>     }
>     private static String amq() {
>         String amq = "activemq:data";
>         amq += "?transacted=true";
>         return amq;
>     }
> }
> {code}



--
This message was sent by Atlassian JIRA
(v7.6.3#76005)

Reply via email to