summaryrefslogtreecommitdiff
path: root/java/broker/src
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2008-04-16 11:43:37 +0000
committerRobert Godfrey <rgodfrey@apache.org>2008-04-16 11:43:37 +0000
commit1fdfb841a9787d0f5bacee5489a963aaf522c332 (patch)
treef18b20a6617d78df4bd98f7b26b259ad5ae96117 /java/broker/src
parent48a474cba1f1ecd98a60810c4f02b6bda1e27172 (diff)
downloadqpid-python-1fdfb841a9787d0f5bacee5489a963aaf522c332.tar.gz
QPID-933 : performance tweaks
git-svn-id: https://svn.apache.org/repos/asf/incubator/qpid/branches/M2.1@648672 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java/broker/src')
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMap.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMapImpl.java15
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/output/amqp0_9/ProtocolOutputConverterImpl.java19
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/protocol/AMQPFastProtocolHandler.java4
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/InMemoryMessageHandle.java3
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/SubscriptionImpl.java25
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/security/access/plugins/AllowAll.java8
7 files changed, 32 insertions, 44 deletions
diff --git a/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMap.java b/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMap.java
index b69a917081..5b0f3cf5eb 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMap.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMap.java
@@ -47,7 +47,7 @@ public interface UnacknowledgedMessageMap
void add(long deliveryTag, UnacknowledgedMessage message);
- void collect(long deliveryTag, boolean multiple, List<UnacknowledgedMessage> msgs);
+ void collect(Long deliveryTag, boolean multiple, List<UnacknowledgedMessage> msgs);
boolean contains(long deliveryTag) throws AMQException;
diff --git a/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMapImpl.java b/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMapImpl.java
index 20ee646a40..5204f13e81 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMapImpl.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/ack/UnacknowledgedMessageMapImpl.java
@@ -28,9 +28,6 @@ import java.util.Map;
import java.util.Set;
import org.apache.qpid.AMQException;
-import org.apache.qpid.framing.AMQShortString;
-import org.apache.qpid.server.protocol.AMQProtocolSession;
-import org.apache.qpid.server.queue.AMQMessage;
import org.apache.qpid.server.txn.TransactionalContext;
public class UnacknowledgedMessageMapImpl implements UnacknowledgedMessageMap
@@ -51,13 +48,7 @@ public class UnacknowledgedMessageMapImpl implements UnacknowledgedMessageMap
_map = new LinkedHashMap<Long, UnacknowledgedMessage>(prefetchLimit);
}
- /*public UnacknowledgedMessageMapImpl(Object lock, Map<Long, UnacknowledgedMessage> map)
- {
- _lock = lock;
- _map = map;
- } */
-
- public void collect(long deliveryTag, boolean multiple, List<UnacknowledgedMessage> msgs)
+ public void collect(Long deliveryTag, boolean multiple, List<UnacknowledgedMessage> msgs)
{
if (multiple)
{
@@ -213,14 +204,14 @@ public class UnacknowledgedMessageMapImpl implements UnacknowledgedMessageMap
}
}
- private void collect(long key, List<UnacknowledgedMessage> msgs)
+ private void collect(Long key, List<UnacknowledgedMessage> msgs)
{
synchronized (_lock)
{
for (Map.Entry<Long, UnacknowledgedMessage> entry : _map.entrySet())
{
msgs.add(entry.getValue());
- if (entry.getKey() == key)
+ if (entry.getKey().equals(key))
{
break;
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_9/ProtocolOutputConverterImpl.java b/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_9/ProtocolOutputConverterImpl.java
index 4bc53dfe03..48d2ca9bc9 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_9/ProtocolOutputConverterImpl.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/output/amqp0_9/ProtocolOutputConverterImpl.java
@@ -179,26 +179,31 @@ public class ProtocolOutputConverterImpl implements ProtocolOutputConverter
private AMQBody createEncodedDeliverFrame(AMQMessage message, final int channelId, final long deliveryTag, final AMQShortString consumerTag)
throws AMQException
{
+
+
final MessagePublishInfo pb = message.getMessagePublishInfo();
final AMQMessageHandle messageHandle = message.getMessageHandle();
- final boolean isRedelivered = messageHandle.isRedelivered();
- final AMQShortString exchangeName = pb.getExchange();
- final AMQShortString routingKey = pb.getRoutingKey();
-
final AMQBody returnBlock = new AMQBody()
{
+
+
+ private final boolean _isRedelivered = messageHandle.isRedelivered();
+ private final AMQShortString _exchangeName = pb.getExchange();
+ private final AMQShortString _routingKey = pb.getRoutingKey();
+
+
public AMQBody _underlyingBody;
public AMQBody createAMQBody()
{
return METHOD_REGISTRY.createBasicDeliverBody(consumerTag,
deliveryTag,
- isRedelivered,
- exchangeName,
- routingKey);
+ _isRedelivered,
+ _exchangeName,
+ _routingKey);
diff --git a/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQPFastProtocolHandler.java b/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQPFastProtocolHandler.java
index d8dbf97e49..ad1c507c04 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQPFastProtocolHandler.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/protocol/AMQPFastProtocolHandler.java
@@ -265,10 +265,6 @@ public class AMQPFastProtocolHandler extends IoHandlerAdapter
*/
public void messageSent(IoSession protocolSession, Object object) throws Exception
{
- if (_logger.isDebugEnabled())
- {
- _logger.debug("Message sent: " + object);
- }
}
protected boolean isSSLClient(ConnectorConfiguration connectionConfig,
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/InMemoryMessageHandle.java b/java/broker/src/main/java/org/apache/qpid/server/queue/InMemoryMessageHandle.java
index 630186991b..0b40f01f1a 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/queue/InMemoryMessageHandle.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/queue/InMemoryMessageHandle.java
@@ -22,6 +22,7 @@ package org.apache.qpid.server.queue;
import java.util.LinkedList;
import java.util.List;
+import java.util.ArrayList;
import org.apache.qpid.AMQException;
import org.apache.qpid.framing.BasicContentHeaderProperties;
@@ -40,7 +41,7 @@ public class InMemoryMessageHandle implements AMQMessageHandle
private MessagePublishInfo _messagePublishInfo;
- private List<ContentChunk> _contentBodies = new LinkedList<ContentChunk>();
+ private List<ContentChunk> _contentBodies = new ArrayList<ContentChunk>();
private boolean _redelivered;
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 bde3ad8ec9..05cd461582 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
@@ -254,13 +254,6 @@ public class SubscriptionImpl implements Subscription
{
long deliveryTag = channel.getNextDeliveryTag();
- // We don't need to add the message to the unacknowledgedMap as we don't need to know if the client
- // received the message. If it is lost in transit that is not important.
-// if (_acks)
-// {
-// channel.addUnacknowledgedBrowsedMessage(msg, deliveryTag, consumerTag, queue);
-// }
-
if (_sendLock.get())
{
_logger.error("Sending " + msg + " when subscriber(" + this + ") is closed!");
@@ -283,25 +276,23 @@ public class SubscriptionImpl implements Subscription
// The send may of course still fail, in which case, as
// the message is unacked, it will be lost.
+ final AMQMessage message = entry.getMessage();
+
if (!_acks)
{
if (_logger.isDebugEnabled())
{
- _logger.debug("No ack mode so dequeuing message immediately: " + entry.getMessage().getMessageId());
+ _logger.debug("No ack mode so dequeuing message immediately: " + message.getMessageId());
}
queue.dequeue(storeContext, entry);
}
-/*
- if (_sendLock.get())
- {
- _logger.error("Sending " + entry + " when subscriber(" + this + ") is closed!");
- }
-*/
+ final ProtocolOutputConverter outputConverter = protocolSession.getProtocolOutputConverter();
+ final int channelId = channel.getChannelId();
synchronized (channel)
{
- long deliveryTag = channel.getNextDeliveryTag();
+ final long deliveryTag = channel.getNextDeliveryTag();
if (_acks)
@@ -309,13 +300,13 @@ public class SubscriptionImpl implements Subscription
channel.addUnacknowledgedMessage(entry, deliveryTag, consumerTag);
}
- protocolSession.getProtocolOutputConverter().writeDeliver(entry.getMessage(), channel.getChannelId(), deliveryTag, consumerTag);
+ outputConverter.writeDeliver(message, channelId, deliveryTag, consumerTag);
}
if (!_acks)
{
- entry.getMessage().decrementReference(storeContext);
+ message.decrementReference(storeContext);
}
}
finally
diff --git a/java/broker/src/main/java/org/apache/qpid/server/security/access/plugins/AllowAll.java b/java/broker/src/main/java/org/apache/qpid/server/security/access/plugins/AllowAll.java
index 9b784069dd..dee1676632 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/security/access/plugins/AllowAll.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/security/access/plugins/AllowAll.java
@@ -28,14 +28,18 @@ import org.apache.qpid.server.security.access.AccessResult;
import org.apache.qpid.server.security.access.Accessable;
import org.apache.qpid.server.security.access.Permission;
import org.apache.commons.configuration.Configuration;
+import org.apache.log4j.Logger;
public class AllowAll implements ACLPlugin
{
+
+ private static final Logger _logger = ACLManager.getLogger();
+
public AccessResult authorise(AMQProtocolSession session, Permission permission, AMQMethodBody body, Object... parameters)
{
- if (ACLManager.getLogger().isDebugEnabled())
+ if (_logger.isDebugEnabled())
{
- ACLManager.getLogger().debug("Allowing user:" + session.getAuthorizedID() + " for :" + permission.toString()
+ _logger.debug("Allowing user:" + session.getAuthorizedID() + " for :" + permission.toString()
+ " on " + body.getClass().getSimpleName()
+ (parameters == null || parameters.length == 0 ? "" : "-" + accessablesToString(parameters)));
}