INT-4123: Add Prefix to the Key-Value MSs

Fixes spring-projects/spring-integration#2213
JIRA: https://jira.spring.io/browse/INT-4123

Fully different `MessageStore`s can be configured for the same shared
Key-Value data-base.
Since the retrieval logic is based on the keys, that may cause the
unexpected messages expiration via `MessageGroupStoreReaper`.

* To distinguish store instances on the shared store add `prefix`
option to the `AbstractKeyValueMessageStore`

* Deprecate the `GemfireMessageStore` `Cache`-based configuration - `setIgnoreJta()` and `afterPropertiesSet()`.
The `GemfireMessageStore` relies only on an externally configured `Region`.

**Cherry-pick to 4.3.x**

Doc Polishing
This commit is contained in:
Artem Bilan
2017-09-05 18:23:03 -04:00
committed by Gary Russell
parent df55ac95d8
commit 5263ea6dff
8 changed files with 176 additions and 111 deletions

View File

@@ -49,12 +49,56 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@Deprecated
protected static final String CREATED_DATE = "CREATED_DATE";
private final String messagePrefix;
private final String groupPrefix;
protected AbstractKeyValueMessageStore() {
this("");
}
/**
* Construct an instance based on the provided prefix for keys to distinguish between
* different store instances in the same target key-value data base. Defaults to an
* empty string - no prefix. The actual prefix for messages is
* {@code prefix + MESSAGE_}; for message groups - {@code prefix + MESSAGE_GROUP_}
* @param prefix the prefix to use
* @since 4.3.12
*/
protected AbstractKeyValueMessageStore(String prefix) {
Assert.notNull(prefix, "'prefix' must not be null");
this.messagePrefix = prefix + MESSAGE_KEY_PREFIX;
this.groupPrefix = prefix + MESSAGE_GROUP_KEY_PREFIX;
}
/**
* Return the configured prefix for message keys to distinguish between different
* store instances in the same target key-value data base. Defaults to the
* {@value MESSAGE_KEY_PREFIX} - without a custom prefix.
* @return the prefix for keys
* @since 4.3.12
*/
protected String getMessagePrefix() {
return this.messagePrefix;
}
/**
* Return the configured prefix for message group keys to distinguish between
* different store instances in the same target key-value data base. Defaults to the
* {@value MESSAGE_GROUP_KEY_PREFIX} - without custom prefix.
* @return the prefix for keys
* @since 4.3.12
*/
public String getGroupPrefix() {
return this.groupPrefix;
}
// MessageStore methods
@Override
public Message<?> getMessage(UUID messageId) {
Assert.notNull(messageId, "'messageId' must not be null");
Object object = doRetrieve(MESSAGE_KEY_PREFIX + messageId);
Object object = doRetrieve(this.messagePrefix + messageId);
if (object != null) {
return extractMessage(object);
}
@@ -80,7 +124,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@Override
public MessageMetadata getMessageMetadata(UUID messageId) {
Assert.notNull(messageId, "'messageId' must not be null");
Object object = doRetrieve(MESSAGE_KEY_PREFIX + messageId);
Object object = doRetrieve(this.messagePrefix + messageId);
if (object != null) {
extractMessage(object);
if (object instanceof MessageHolder) {
@@ -100,13 +144,13 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
protected void doAddMessage(Message<?> message) {
Assert.notNull(message, "'message' must not be null");
UUID messageId = message.getHeaders().getId();
doStoreIfAbsent(MESSAGE_KEY_PREFIX + messageId, new MessageHolder(message));
doStoreIfAbsent(this.messagePrefix + messageId, new MessageHolder(message));
}
@Override
public Message<?> removeMessage(UUID id) {
Assert.notNull(id, "'id' must not be null");
Object object = doRemove(MESSAGE_KEY_PREFIX + id);
Object object = doRemove(this.messagePrefix + id);
if (object != null) {
return extractMessage(object);
}
@@ -118,7 +162,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@Override
@ManagedAttribute
public long getMessageCount() {
Collection<?> messageIds = doListKeys(MESSAGE_KEY_PREFIX + "*");
Collection<?> messageIds = doListKeys(this.messagePrefix + "*");
return (messageIds != null) ? messageIds.size() : 0;
}
@@ -147,7 +191,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@Override
public MessageGroupMetadata getGroupMetadata(Object groupId) {
Assert.notNull(groupId, "'groupId' must not be null");
Object mgm = this.doRetrieve(MESSAGE_GROUP_KEY_PREFIX + groupId);
Object mgm = this.doRetrieve(this.groupPrefix + groupId);
if (mgm != null) {
Assert.isInstanceOf(MessageGroupMetadata.class, mgm);
return (MessageGroupMetadata) mgm;
@@ -187,7 +231,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
// store MessageGroupMetadata built from enriched MG
doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, metadata);
doStore(this.groupPrefix + groupId, metadata);
}
@Override
@@ -195,17 +239,17 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
Assert.notNull(groupId, "'groupId' must not be null");
Assert.notNull(messages, "'messages' must not be null");
Object mgm = doRetrieve(MESSAGE_GROUP_KEY_PREFIX + groupId);
Object mgm = doRetrieve(this.groupPrefix + groupId);
if (mgm != null) {
Assert.isInstanceOf(MessageGroupMetadata.class, mgm);
MessageGroupMetadata messageGroupMetadata = (MessageGroupMetadata) mgm;
for (Message<?> messageToRemove : messages) {
UUID messageId = messageToRemove.getHeaders().getId();
messageGroupMetadata.remove(messageId);
doRemove(MESSAGE_KEY_PREFIX + messageId);
doRemove(this.messagePrefix + messageId);
}
messageGroupMetadata.setLastModified(System.currentTimeMillis());
doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, messageGroupMetadata);
doStore(this.groupPrefix + groupId, messageGroupMetadata);
}
}
@@ -216,7 +260,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
if (metadata != null) {
metadata.complete();
metadata.setLastModified(System.currentTimeMillis());
doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, metadata);
doStore(this.groupPrefix + groupId, metadata);
}
}
@@ -226,7 +270,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@Override
public void removeMessageGroup(Object groupId) {
Assert.notNull(groupId, "'groupId' must not be null");
Object mgm = doRemove(MESSAGE_GROUP_KEY_PREFIX + groupId);
Object mgm = doRemove(this.groupPrefix + groupId);
if (mgm != null) {
Assert.isInstanceOf(MessageGroupMetadata.class, mgm);
MessageGroupMetadata messageGroupMetadata = (MessageGroupMetadata) mgm;
@@ -248,7 +292,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
}
metadata.setLastReleasedMessageSequenceNumber(sequenceNumber);
metadata.setLastModified(System.currentTimeMillis());
doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, metadata);
doStore(this.groupPrefix + groupId, metadata);
}
@Override
@@ -259,7 +303,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
if (firstId != null) {
groupMetadata.remove(firstId);
groupMetadata.setLastModified(System.currentTimeMillis());
doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, groupMetadata);
doStore(this.groupPrefix + groupId, groupMetadata);
return removeMessage(firstId);
}
}
@@ -295,7 +339,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
@SuppressWarnings("unchecked")
public Iterator<MessageGroup> iterator() {
final Iterator<?> idIterator = normalizeKeys(
(Collection<String>) doListKeys(MESSAGE_GROUP_KEY_PREFIX + "*"))
(Collection<String>) doListKeys(this.groupPrefix + "*"))
.iterator();
return new MessageGroupIterator(idIterator);
}
@@ -304,11 +348,11 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
Set<String> normalizedKeys = new HashSet<String>();
for (Object key : keys) {
String strKey = (String) key;
if (strKey.startsWith(MESSAGE_GROUP_KEY_PREFIX)) {
strKey = strKey.replace(MESSAGE_GROUP_KEY_PREFIX, "");
if (strKey.startsWith(this.groupPrefix)) {
strKey = strKey.replace(this.groupPrefix, "");
}
else if (strKey.startsWith(MESSAGE_KEY_PREFIX)) {
strKey = strKey.replace(MESSAGE_KEY_PREFIX, "");
else if (strKey.startsWith(this.messagePrefix)) {
strKey = strKey.replace(this.messagePrefix, "");
}
normalizedKeys.add(strKey);
}
@@ -359,6 +403,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
public void remove() {
throw new UnsupportedOperationException();
}
}
}