You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@camel.apache.org by bv...@apache.org on 2012/02/26 16:30:00 UTC
svn commit: r1293852 - in /camel/trunk/camel-core/src:
main/java/org/apache/camel/api/management/mbean/
main/java/org/apache/camel/management/
main/java/org/apache/camel/management/mbean/
main/java/org/apache/camel/processor/idempotent/ test/java/org/a...
Author: bvahdat
Date: Sun Feb 26 15:29:59 2012
New Revision: 1293852
URL: http://svn.apache.org/viewvc?rev=1293852&view=rev
Log:
CAMEL-4782: Add ManagedIdempotentConsumerBean providing get & reset of the duplicate messages count detected.
Added:
camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java (with props)
camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java (with props)
Modified:
camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
Added: camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java?rev=1293852&view=auto
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java (added)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java Sun Feb 26 15:29:59 2012
@@ -0,0 +1,30 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.api.management.mbean;
+
+import org.apache.camel.api.management.ManagedAttribute;
+import org.apache.camel.api.management.ManagedOperation;
+
+public interface ManagedIdempotentConsumerMBean extends ManagedProcessorMBean {
+
+ @ManagedAttribute(description = "Current count of duplicate Messages")
+ long getDuplicateMessageCount();
+
+ @ManagedOperation(description = "Reset the current count of duplicate Messages")
+ void resetDuplicateMessageCount();
+
+}
Propchange: camel/trunk/camel-core/src/main/java/org/apache/camel/api/management/mbean/ManagedIdempotentConsumerMBean.java
------------------------------------------------------------------------------
svn:eol-style = native
Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementLifecycleStrategy.java Sun Feb 26 15:29:59 2012
@@ -422,10 +422,9 @@ public class DefaultManagementLifecycleS
ManagedService ms = (ManagedService) answer;
ms.setRoute(route);
ms.init(getManagementStrategy());
- return answer;
- } else {
- return answer;
}
+
+ return answer;
}
private Object getManagedObjectForProcessor(CamelContext context, Processor processor, Route route) {
Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/DefaultManagementObjectStrategy.java Sun Feb 26 15:29:59 2012
@@ -39,6 +39,7 @@ import org.apache.camel.management.mbean
import org.apache.camel.management.mbean.ManagedEndpoint;
import org.apache.camel.management.mbean.ManagedErrorHandler;
import org.apache.camel.management.mbean.ManagedEventNotifier;
+import org.apache.camel.management.mbean.ManagedIdempotentConsumer;
import org.apache.camel.management.mbean.ManagedProcessor;
import org.apache.camel.management.mbean.ManagedProducer;
import org.apache.camel.management.mbean.ManagedRoute;
@@ -54,6 +55,7 @@ import org.apache.camel.processor.Delaye
import org.apache.camel.processor.ErrorHandler;
import org.apache.camel.processor.SendProcessor;
import org.apache.camel.processor.Throttler;
+import org.apache.camel.processor.idempotent.IdempotentConsumer;
import org.apache.camel.spi.BrowsableEndpoint;
import org.apache.camel.spi.EventNotifier;
import org.apache.camel.spi.ManagementObjectStrategy;
@@ -178,6 +180,8 @@ public class DefaultManagementObjectStra
answer = new ManagedSendProcessor(context, (SendProcessor) target, definition);
} else if (target instanceof BeanProcessor) {
answer = new ManagedBeanProcessor(context, (BeanProcessor) target, definition);
+ } else if (target instanceof IdempotentConsumer) {
+ answer = new ManagedIdempotentConsumer(context, (IdempotentConsumer) target, definition);
} else if (target instanceof org.apache.camel.spi.ManagementAware) {
return ((org.apache.camel.spi.ManagementAware<Processor>) target).getManagedObject(processor);
}
Added: camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java?rev=1293852&view=auto
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java (added)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java Sun Feb 26 15:29:59 2012
@@ -0,0 +1,47 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.management.mbean;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.api.management.ManagedResource;
+import org.apache.camel.api.management.mbean.ManagedIdempotentConsumerMBean;
+import org.apache.camel.model.ProcessorDefinition;
+import org.apache.camel.processor.idempotent.IdempotentConsumer;
+
+@ManagedResource(description = "Managed Idempotent Consumer")
+public class ManagedIdempotentConsumer extends ManagedProcessor implements ManagedIdempotentConsumerMBean {
+
+ public ManagedIdempotentConsumer(CamelContext context, IdempotentConsumer idempotentConsumer, ProcessorDefinition<?> definition) {
+ super(context, idempotentConsumer, definition);
+ }
+
+ @Override
+ public IdempotentConsumer getProcessor() {
+ return (IdempotentConsumer) super.getProcessor();
+ }
+
+ @Override
+ public long getDuplicateMessageCount() {
+ return getProcessor().getDuplicateMessageCount();
+ }
+
+ @Override
+ public void resetDuplicateMessageCount() {
+ getProcessor().resetDuplicateMessageCount();
+ }
+
+}
Propchange: camel/trunk/camel-core/src/main/java/org/apache/camel/management/mbean/ManagedIdempotentConsumer.java
------------------------------------------------------------------------------
svn:eol-style = native
Modified: camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java (original)
+++ camel/trunk/camel-core/src/main/java/org/apache/camel/processor/idempotent/IdempotentConsumer.java Sun Feb 26 15:29:59 2012
@@ -18,6 +18,7 @@ package org.apache.camel.processor.idemp
import java.util.ArrayList;
import java.util.List;
+import java.util.concurrent.atomic.AtomicLong;
import org.apache.camel.AsyncCallback;
import org.apache.camel.AsyncProcessor;
@@ -45,6 +46,7 @@ public class IdempotentConsumer extends
private final boolean eager;
private final boolean skipDuplicate;
private final boolean removeOnFailure;
+ private final AtomicLong duplicateMessageCount = new AtomicLong();
public IdempotentConsumer(Expression messageIdExpression, IdempotentRepository<String> idempotentRepository,
boolean eager, boolean skipDuplicate, boolean removeOnFailure, Processor processor) {
@@ -84,11 +86,10 @@ public class IdempotentConsumer extends
if (!newKey) {
// mark the exchange as duplicate
exchange.setProperty(Exchange.DUPLICATE_MESSAGE, Boolean.TRUE);
- }
- if (!newKey) {
// we already have this key so its a duplicate message
- onDuplicateMessage(exchange, messageId);
+ onDuplicate(exchange, messageId);
+
if (skipDuplicate) {
// if we should skip duplicate then we are done
LOG.debug("Ignoring duplicate message with id: {} for exchange: {}", messageId, exchange);
@@ -131,6 +132,10 @@ public class IdempotentConsumer extends
return processor;
}
+ public long getDuplicateMessageCount() {
+ return duplicateMessageCount.get();
+ }
+
// Implementation methods
// -------------------------------------------------------------------------
@@ -143,6 +148,19 @@ public class IdempotentConsumer extends
}
/**
+ * Resets the duplicate message counter to <code>0L</code>.
+ */
+ public void resetDuplicateMessageCount() {
+ duplicateMessageCount.set(0L);
+ }
+
+ private void onDuplicate(Exchange exchange, String messageId) {
+ duplicateMessageCount.incrementAndGet();
+
+ onDuplicateMessage(exchange, messageId);
+ }
+
+ /**
* A strategy method to allow derived classes to overload the behaviour of
* processing a duplicate message
*
Modified: camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java
URL: http://svn.apache.org/viewvc/camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java?rev=1293852&r1=1293851&r2=1293852&view=diff
==============================================================================
--- camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java (original)
+++ camel/trunk/camel-core/src/test/java/org/apache/camel/management/ManagedMemoryIdempotentConsumerTest.java Sun Feb 26 15:29:59 2012
@@ -28,7 +28,6 @@ import org.apache.camel.builder.RouteBui
import org.apache.camel.component.mock.MockEndpoint;
import org.apache.camel.processor.idempotent.MemoryIdempotentRepository;
import org.apache.camel.spi.IdempotentRepository;
-import org.apache.camel.util.CastUtils;
/**
* @version
@@ -42,7 +41,7 @@ public class ManagedMemoryIdempotentCons
MBeanServer mbeanServer = getMBeanServer();
// services
- Set<ObjectName> names = CastUtils.cast(mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=services,*"), null));
+ Set<ObjectName> names = mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=services,*"), null);
ObjectName on = null;
for (ObjectName name : names) {
if (name.toString().contains("MemoryIdempotentRepository")) {
@@ -93,6 +92,57 @@ public class ManagedMemoryIdempotentCons
assertTrue(repo.contains("4"));
}
+ public void testDuplicateMessagesCountAreCorrectlyCounted() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+
+ // processors
+ Set<ObjectName> names = mbeanServer.queryNames(new ObjectName("org.apache.camel" + ":type=processors,*"), null);
+ ObjectName on = null;
+ for (ObjectName name : names) {
+ if (name.toString().contains("idempotentConsumer")) {
+ on = name;
+ break;
+ }
+ }
+ assertTrue("Should be registered", mbeanServer.isRegistered(on));
+
+ Long count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+ assertEquals(0L, count.longValue());
+
+ resultEndpoint.expectedBodiesReceived("one", "two");
+
+ sendMessage("1", "one");
+ sendMessage("2", "two");
+ sendMessage("1", "one");
+ sendMessage("2", "two");
+
+ resultEndpoint.assertIsSatisfied();
+
+ count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+ assertEquals(2L, count.longValue());
+
+ // reset the count
+ mbeanServer.invoke(on, "resetDuplicateMessageCount", null, null);
+
+ // count should be resetted
+ count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+ assertEquals(0L, count.longValue());
+
+ resetMocks();
+
+ resultEndpoint.expectedBodiesReceived("five");
+
+ sendMessage("4", "four");
+ sendMessage("4", "four");
+ sendMessage("5", "five");
+ sendMessage("4", "four");
+
+ resultEndpoint.assertIsSatisfied();
+
+ count = (Long) mbeanServer.getAttribute(on, "DuplicateMessageCount");
+ assertEquals(3L, count.longValue());
+ }
+
protected void sendMessage(final Object messageId, final Object body) {
template.send(startEndpoint, new Processor() {
public void process(Exchange exchange) {