summaryrefslogtreecommitdiff
path: root/qpid/java/broker-plugins/amqp-0-8-protocol
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2014-11-04 15:52:59 +0000
committerRobert Godfrey <rgodfrey@apache.org>2014-11-04 15:52:59 +0000
commitc64d0182543fd9b2b8029fb18f99993a3891977c (patch)
treeb77207767e0fe87c6b9cdde0272d85a95f2ae5f4 /qpid/java/broker-plugins/amqp-0-8-protocol
parent930ffe78c89a3937b38c61bc5c7b324e701e05f6 (diff)
downloadqpid-python-c64d0182543fd9b2b8029fb18f99993a3891977c.tar.gz
QPID-6207 : [Java Broker] Flow uncommitted messages to disk if combined size greater than threshold
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1636617 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-plugins/amqp-0-8-protocol')
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java49
1 files changed, 49 insertions, 0 deletions
diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
index 37ac1f84c4..012d7bffd6 100644
--- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
+++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
@@ -82,6 +82,7 @@ import org.apache.qpid.server.message.ServerMessage;
import org.apache.qpid.server.model.Broker;
import org.apache.qpid.server.model.ConfigurationChangeListener;
import org.apache.qpid.server.model.ConfiguredObject;
+import org.apache.qpid.server.model.Connection;
import org.apache.qpid.server.model.Consumer;
import org.apache.qpid.server.model.Exchange;
import org.apache.qpid.server.model.ExclusivityPolicy;
@@ -206,6 +207,9 @@ public class AMQChannel
private long _blockingTimeout;
private boolean _confirmOnPublish;
private long _confirmedMessageCounter;
+ private volatile long _uncommittedMessageSize;
+ private final List<StoredMessage<MessageMetaData>> _uncommittedMessages = new ArrayList<>();
+ private long _maxUncommittedInMemorySize;
public AMQChannel(AMQProtocolEngine connection, int channelId, final MessageStore messageStore)
{
@@ -216,6 +220,7 @@ public class AMQChannel
connection.getAuthorizedSubject().getPublicCredentials(),
connection.getAuthorizedSubject().getPrivateCredentials());
_subject.getPrincipals().add(new SessionPrincipal(this));
+ _maxUncommittedInMemorySize = connection.getVirtualHost().getContextValue(Long.class, Connection.MAX_UNCOMMITTED_IN_MEMORY_SIZE);
_logSubject = new ChannelLogSubject(this);
_messageStore = messageStore;
@@ -481,6 +486,7 @@ public class AMQChannel
.createBasicAckBody(_confirmedMessageCounter, false);
_connection.writeFrame(responseBody.generateFrame(_channelId));
}
+ incrementUncommittedMessageSize(handle);
incrementOutstandingTxnsIfNecessary();
}
}
@@ -506,6 +512,36 @@ public class AMQChannel
}
+ private void incrementUncommittedMessageSize(final StoredMessage<MessageMetaData> handle)
+ {
+ if (isTransactional())
+ {
+ _uncommittedMessageSize += handle.getMetaData().getContentSize();
+ if (_uncommittedMessageSize > getMaxUncommittedInMemorySize())
+ {
+ handle.flowToDisk();
+ if(!_uncommittedMessages.isEmpty() || _uncommittedMessageSize == handle.getMetaData().getContentSize())
+ {
+ getVirtualHost().getEventLogger()
+ .message(_logSubject, ChannelMessages.LARGE_TRANSACTION_WARN(_uncommittedMessageSize));
+ }
+
+ if(!_uncommittedMessages.isEmpty())
+ {
+ for (StoredMessage<MessageMetaData> uncommittedHandle : _uncommittedMessages)
+ {
+ uncommittedHandle.flowToDisk();
+ }
+ _uncommittedMessages.clear();
+ }
+ }
+ else
+ {
+ _uncommittedMessages.add(handle);
+ }
+ }
+ }
+
/**
* Either throws a {@link AMQConnectionException} or returns the message
*
@@ -1182,6 +1218,13 @@ public class AMQChannel
_txnStarts.incrementAndGet();
decrementOutstandingTxnsIfNecessary();
}
+ resetUncommittedMessages();
+ }
+
+ private void resetUncommittedMessages()
+ {
+ _uncommittedMessageSize = 0l;
+ _uncommittedMessages.clear();
}
public void rollback(Runnable postRollbackTask)
@@ -1209,6 +1252,7 @@ public class AMQChannel
_txnRejects.incrementAndGet();
_txnStarts.incrementAndGet();
decrementOutstandingTxnsIfNecessary();
+ resetUncommittedMessages();
}
postRollbackTask.run();
@@ -1368,6 +1412,11 @@ public class AMQChannel
return _currentMessage != null;
}
+ public long getMaxUncommittedInMemorySize()
+ {
+ return _maxUncommittedInMemorySize;
+ }
+
private class GetDeliveryMethod implements ClientDeliveryMethod
{