I'm using Camel 2.8.5, and am finding that messages sent to SEDA endpoints
that have polling consumers are sometimes lost - when sending more than 1000
messages.

My example scenario is shown below.  I have a timer endpoint to kick off my
experiment.  The route then proceeds as follows:

- Generate a list of 2000 Strings out of nowhere
- Split the message into 2000 separate messages
- For each message:
-- inOnly to SEDA queue 1
-- inOnly to SEDA queue 2
-- end
- use a ConsumerTemplate to poll all messages from SEDA queue 1 until no
more messages are found; count these messages
- display the total count of messages for queue 1

Separately, I have another route to count the messages to queue 2 in an
event-driven manner.  The count for this queue is shown alongside the count
for queue 1.

The ConsumerTemplate always finds at least 1000 messages on queue 1, and
sometimes more.  But it never finds all 2000.  Typical is very close to
1000.  The event-driven counter always finds all 2000.  Why does the polling
method not find all messages?

Some extra info:
- if I start with more than 2000 messages, the polling consumer still
generally finds around 1000 messages; the event-driven consumer finds them
all still
- if I start with less than 1000 messages, both consumers find all their
messages
- if I don't define a 'size' parameter on my polling consumer's queue, the
behaviour does not change

>>>>

package com.example;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;

import org.apache.camel.ConsumerTemplate;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.builder.RouteBuilder;
import org.springframework.beans.factory.annotation.Autowired;

public class ConsumerTestRouteBuilder extends RouteBuilder {

    private static final String EVENT_DRIVEN_Q = "seda:eventDrivenQ";
    private static final String POLLING_Q = "seda:pollingQ?size=10000";

 @Autowired
    private ConsumerTemplate consumer;
 
 private AtomicInteger eventDrivenCount = new AtomicInteger();
 private AtomicInteger pollingCount = new AtomicInteger();

 @Override
 public void configure() throws Exception {
  
     from("timer://foo?fixedRate=true&period=20000")
         .process(new Processor() {
                @Override
                public void process(Exchange exchange) throws Exception {
                    System.out.println("Start of test");
                    eventDrivenCount.set(0);
                    pollingCount.set(0);
                    List<String> answer = new ArrayList<String>();
                    for(int i=0;i<2000;i++) {
                        answer.add("myString"+i);
                    }
                    exchange.getIn().setBody(answer);
                }
            })
            .split(body())
                .inOnly(POLLING_Q)
                .inOnly(EVENT_DRIVEN_Q)
                .end()
            .process(new Processor() {
                @Override
                public void process(Exchange exchange) throws Exception {
                    while(consumer.receive(POLLING_Q, 7000) !=null ) {
                        pollingCount.incrementAndGet();
                    }
                };
            })
            .process(new Processor() {
                @Override
                public void process(Exchange exchange) throws Exception {
                    System.out.println("Event-driven consumptions =
"+eventDrivenCount.get());
                    System.out.println("Polling consumptions =
"+pollingCount.get());
                    System.out.println("End of test");
                    System.out.println();
                }
            });
     
     from(EVENT_DRIVEN_Q)
         .process(new Processor() {
                @Override
                public void process(Exchange exchange) throws Exception {
                    eventDrivenCount.incrementAndGet();
                }
         });
 } 
}

<<<<
I configure using Spring:
>>>>

<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans";
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance";
 xmlns:cxf="http://camel.apache.org/schema/cxf";
xmlns:camel="http://camel.apache.org/schema/spring";
xmlns:context="http://www.springframework.org/schema/context";
 xmlns:jaxws="http://cxf.apache.org/jaxws";
 xsi:schemaLocation="http://www.springframework.org/schema/beans 
                   
http://www.springframework.org/schema/beans/spring-beans.xsd
                    http://camel.apache.org/schema/spring
                    http://camel.apache.org/schema/spring/camel-spring.xsd
                    http://www.springframework.org/schema/context
                   
http://www.springframework.org/schema/context/spring-context.xsd";>

 <bean id="consumerTestRouteBuilder"
class="com.example.ConsumerTestRouteBuilder" />

    <context:annotation-config />

 <camel:camelContext id="camelContext"
xmlns="http://camel.apache.org/schema/spring";>
        <camel:consumerTemplate id="consumer" />
  <camel:routeBuilder ref="consumerTestRouteBuilder" />
 </camel:camelContext>

</beans>

<<<<




--
View this message in context: 
http://camel.465427.n5.nabble.com/Losing-Messages-in-SEDA-endpoint-with-polling-consumer-tp5717711.html
Sent from the Camel - Users mailing list archive at Nabble.com.

Reply via email to