diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-11-04 15:52:59 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-11-04 15:52:59 +0000 |
| commit | c64d0182543fd9b2b8029fb18f99993a3891977c (patch) | |
| tree | b77207767e0fe87c6b9cdde0272d85a95f2ae5f4 /qpid/java/broker-plugins/amqp-0-8-protocol | |
| parent | 930ffe78c89a3937b38c61bc5c7b324e701e05f6 (diff) | |
| download | qpid-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.java | 49 |
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 { |
