diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2008-04-16 11:43:37 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2008-04-16 11:43:37 +0000 |
| commit | 1fdfb841a9787d0f5bacee5489a963aaf522c332 (patch) | |
| tree | f18b20a6617d78df4bd98f7b26b259ad5ae96117 /java/broker/src | |
| parent | 48a474cba1f1ecd98a60810c4f02b6bda1e27172 (diff) | |
| download | qpid-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')
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))); } |
