summaryrefslogtreecommitdiff
path: root/java/broker
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2007-05-09 10:24:13 +0000
committerRobert Godfrey <rgodfrey@apache.org>2007-05-09 10:24:13 +0000
commitcf41d7497d9fbd90cc772396bf3b749d959c1e16 (patch)
tree9f94367bae103343cacaea9945c567df5dc5d318 /java/broker
parente4557cabd08f89240d6100479604062d4e9ae31e (diff)
downloadqpid-python-cf41d7497d9fbd90cc772396bf3b749d959c1e16.tar.gz
QPID-482 : Small performance enhancements
git-svn-id: https://svn.apache.org/repos/asf/incubator/qpid/branches/M2@536480 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java/broker')
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/AMQChannel.java8
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/output/amqp0_8/ProtocolOutputConverterImpl.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/protocol/AMQMinaProtocolSession.java9
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java30
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java58
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/state/AMQStateManager.java4
6 files changed, 52 insertions, 59 deletions
diff --git a/java/broker/src/main/java/org/apache/qpid/server/AMQChannel.java b/java/broker/src/main/java/org/apache/qpid/server/AMQChannel.java
index b48953ad33..33c87d90e3 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/AMQChannel.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/AMQChannel.java
@@ -115,6 +115,8 @@ public class AMQChannel
//Why do we need this reference ? - ritchiem
private final AMQProtocolSession _session;
+ private boolean _closing;
+
public AMQChannel(AMQProtocolSession session, int channelId, MessageStore messageStore, MessageRouter exchanges)
throws AMQException
@@ -353,6 +355,7 @@ public class AMQChannel
*/
public void close(AMQProtocolSession session) throws AMQException
{
+ _closing = true;
_txnContext.rollback();
unsubscribeAllConsumers(session);
requeue();
@@ -961,4 +964,9 @@ public class AMQChannel
{
return _txnContext;
}
+
+ public boolean isClosing()
+ {
+ return _closing;
+ }
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_8/ProtocolOutputConverterImpl.java b/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_8/ProtocolOutputConverterImpl.java
index bd5bb632fe..28b2489142 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_8/ProtocolOutputConverterImpl.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_8/ProtocolOutputConverterImpl.java
@@ -77,7 +77,7 @@ public class ProtocolOutputConverterImpl implements ProtocolOutputConverter
final AMQMessageHandle messageHandle = message.getMessageHandle();
final StoreContext storeContext = message.getStoreContext();
- final long messageId = message.getMessageId();
+ final Long messageId = message.getMessageId();
final int bodyCount = messageHandle.getBodyCount(storeContext,messageId);
diff --git a/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQMinaProtocolSession.java b/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQMinaProtocolSession.java
index 3162efa183..d430f1af94 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQMinaProtocolSession.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQMinaProtocolSession.java
@@ -432,15 +432,16 @@ public class AMQMinaProtocolSession implements AMQProtocolSession,
public AMQChannel getChannel(int channelId) throws AMQException
{
- if (channelAwaitingClosure(channelId))
+ final AMQChannel channel = ((channelId & CHANNEL_CACHE_SIZE) == channelId)
+ ? _cachedChannels[channelId]
+ : _channelMap.get(channelId);
+ if (channel == null || channel.isClosing())
{
return null;
}
else
{
- return ((channelId & CHANNEL_CACHE_SIZE) == channelId)
- ? _cachedChannels[channelId]
- : _channelMap.get(channelId);
+ return channel;
}
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
index 955aaa6acb..90cf845f10 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
@@ -207,10 +207,6 @@ public class AMQMessage
// _taken = new AtomicBoolean(false);
- if (_log.isDebugEnabled())
- {
- _log.debug("Message(" + System.identityHashCode(this) + ") created (" + debugIdentity() + ")");
- }
}
/**
@@ -355,7 +351,7 @@ public class AMQMessage
return _transientMessageData.isAllContentReceived();
}
- public long getMessageId()
+ public Long getMessageId()
{
return _messageId;
}
@@ -374,10 +370,10 @@ public class AMQMessage
protected void incrementReference()
{
_referenceCount.incrementAndGet();
- if (_log.isDebugEnabled())
- {
- _log.debug("Ref count on message " + debugIdentity() + " incremented " + Arrays.asList(Thread.currentThread().getStackTrace()).subList(3, 6));
- }
+// if (_log.isDebugEnabled())
+// {
+// _log.debug("Ref count on message " + debugIdentity() + " incremented " + Arrays.asList(Thread.currentThread().getStackTrace()).subList(3, 6));
+// }
}
/**
@@ -401,10 +397,10 @@ public class AMQMessage
{
try
{
- if (_log.isDebugEnabled())
- {
- _log.debug("Decremented ref count on message " + debugIdentity() + " is zero; removing message" + Arrays.asList(Thread.currentThread().getStackTrace()).subList(3, 6));
- }
+// if (_log.isDebugEnabled())
+// {
+// _log.debug("Decremented ref count on message " + debugIdentity() + " is zero; removing message" + Arrays.asList(Thread.currentThread().getStackTrace()).subList(3, 6));
+// }
// must check if the handle is null since there may be cases where we decide to throw away a message
// and the handle has not yet been constructed
@@ -422,14 +418,6 @@ public class AMQMessage
}
else
{
- if (_log.isDebugEnabled())
- {
- _log.debug("Decremented ref count is now " + count + " for message id " + debugIdentity() + "\n" + Arrays.asList(Thread.currentThread().getStackTrace()).subList(3, 5));
- if (count < 0)
- {
- Thread.dumpStack();
- }
- }
if (count < 0)
{
throw new MessageCleanupException("Reference count for message id " + debugIdentity() + " has gone below 0.");
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java b/java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java
index 3bce950ba9..c496996002 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java
@@ -322,17 +322,17 @@ public class SubscriptionImpl implements Subscription
public boolean isSuspended()
{
- if (_suspensionlogger.isInfoEnabled())
- {
- if (channel.isSuspended())
- {
- _suspensionlogger.info("Subscription(" + debugIdentity() + ") channel's is susupended");
- }
- if (_sendLock.get())
- {
- _suspensionlogger.info("Subscription(" + debugIdentity() + ") has sendLock set so closing.");
- }
- }
+// if (_suspensionlogger.isInfoEnabled())
+// {
+// if (channel.isSuspended())
+// {
+// _suspensionlogger.debug("Subscription(" + debugIdentity() + ") channel's is susupended");
+// }
+// if (_sendLock.get())
+// {
+// _suspensionlogger.debug("Subscription(" + debugIdentity() + ") has sendLock set so closing.");
+// }
+// }
return channel.isSuspended() || _sendLock.get();
}
@@ -381,11 +381,11 @@ public class SubscriptionImpl implements Subscription
{
if (localInstance == msgInstance || localInstance.equals(msgInstance))
{
- if (_logger.isTraceEnabled())
- {
- _logger.trace("(" + debugIdentity() + ") has no interest as it is a local message(" +
- msg.debugIdentity() + ")");
- }
+// if (_logger.isTraceEnabled())
+// {
+// _logger.trace("(" + debugIdentity() + ") has no interest as it is a local message(" +
+// msg.debugIdentity() + ")");
+// }
return false;
}
}
@@ -399,11 +399,11 @@ public class SubscriptionImpl implements Subscription
msgInstance = publisher.getClientIdentifier();
if (localInstance == msgInstance || ((localInstance != null) && localInstance.equals(msgInstance)))
{
- if (_logger.isTraceEnabled())
- {
- _logger.trace("(" + debugIdentity() + ") has no interest as it is a local message(" +
- msg.debugIdentity() + ")");
- }
+// if (_logger.isTraceEnabled())
+// {
+// _logger.trace("(" + debugIdentity() + ") has no interest as it is a local message(" +
+// msg.debugIdentity() + ")");
+// }
return false;
}
}
@@ -431,18 +431,18 @@ public class SubscriptionImpl implements Subscription
{
if (_filters != null)
{
- if (_logger.isTraceEnabled())
- {
- _logger.trace("(" + debugIdentity() + ") has filters.");
- }
+// if (_logger.isTraceEnabled())
+// {
+// _logger.trace("(" + debugIdentity() + ") has filters.");
+// }
return _filters.allAllow(msg);
}
else
{
- if (_logger.isTraceEnabled())
- {
- _logger.trace("(" + debugIdentity() + ") has no filters");
- }
+// if (_logger.isTraceEnabled())
+// {
+// _logger.trace("(" + debugIdentity() + ") has no filters");
+// }
return true;
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/state/AMQStateManager.java b/java/broker/src/main/java/org/apache/qpid/server/state/AMQStateManager.java
index d12f5cd084..50129ec274 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/state/AMQStateManager.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/state/AMQStateManager.java
@@ -241,10 +241,6 @@ public class AMQStateManager implements AMQMethodListener
B frame)
throws IllegalStateTransitionException
{
- if (_logger.isDebugEnabled())
- {
- _logger.debug("Looking for state transition handler for frame " + frame.getClass());
- }
final Map<Class<? extends AMQMethodBody>, StateAwareMethodListener<? extends AMQMethodBody>>
classToHandlerMap = _state2HandlersMap.get(currentState);