[
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 the doCatch part will
never be reached and that the only way to make a doTry with aggregator (i mean
by direcly call endDoTry). 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()
// Never reached why is valid here /?!\
.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* 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 the doCatch part will
never be reached and that the only way to make a doTry (i mean by direcly call
endDoTry). 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()
// Never reached why is valid here /?!\
.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: Critical
>
> /!\ 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 the doCatch part will
> never be reached and that the only way to make a doTry with aggregator (i
> mean by direcly call endDoTry). 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()
> // Never reached why is valid here /?!\
> .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)