Merge pull request #199 from olegz/INT-2253
added support for updating MessageGroup with lastModified timestamp on every touch
This commit is contained in:
@@ -138,7 +138,8 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
|
||||
}
|
||||
}
|
||||
this.removeMessage(messageToRemove.getHeaders().getId());
|
||||
|
||||
rawGroup.setLastModified(System.currentTimeMillis());
|
||||
|
||||
this.doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, new MessageGroupMetadata(rawGroup));
|
||||
messageGroup = this.getSimpleMessageGroup(this.getMessageGroup(groupId));
|
||||
|
||||
@@ -150,6 +151,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
|
||||
Assert.notNull(groupId, "'groupId' must not be null");
|
||||
SimpleMessageGroup messageGroup = this.buildMessageGroup(groupId, true);
|
||||
messageGroup.complete();
|
||||
messageGroup.setLastModified(System.currentTimeMillis());
|
||||
this.doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, new MessageGroupMetadata(messageGroup));
|
||||
}
|
||||
|
||||
@@ -174,6 +176,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
|
||||
Assert.notNull(groupId, "'groupId' must not be null");
|
||||
SimpleMessageGroup messageGroup = this.buildMessageGroup(groupId, true);
|
||||
messageGroup.setLastReleasedMessageSequenceNumber(sequenceNumber);
|
||||
messageGroup.setLastModified(System.currentTimeMillis());
|
||||
this.doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, new MessageGroupMetadata(messageGroup));
|
||||
}
|
||||
|
||||
@@ -187,6 +190,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS
|
||||
UUID firstId = messageGroupMetadata.firstId();
|
||||
if (firstId != null){
|
||||
messageGroupMetadata.remove(firstId);
|
||||
messageGroupMetadata.setLastModified(System.currentTimeMillis());
|
||||
this.doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, messageGroupMetadata);
|
||||
return this.removeMessage(firstId);
|
||||
}
|
||||
|
||||
@@ -76,6 +76,10 @@ public class MessageGroupMetadata implements Serializable{
|
||||
}
|
||||
this.messageCreationDateToIdMappings.remove(currentTimestamp);
|
||||
}
|
||||
|
||||
public void setLastModified(long lastModified) {
|
||||
this.lastModified = lastModified;
|
||||
}
|
||||
|
||||
public Object getGroupId() {
|
||||
return this.groupId;
|
||||
|
||||
@@ -109,6 +109,8 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
|
||||
private static final String CREATE_MESSAGE_IN_GROUP = "INSERT into %PREFIX%MESSAGE_GROUP(MESSAGE_ID, REGION, CREATED_DATE, UPDATED_DATE, GROUP_KEY, MARKED, COMPLETE, LAST_RELEASED_SEQUENCE)"
|
||||
+ " values (?, ?, ?, ?, ?, 0, 0, 0)";
|
||||
|
||||
private static final String UPDATE_GROUP = "UPDATE %PREFIX%MESSAGE_GROUP set UPDATED_DATE=? where GROUP_KEY=? and REGION=?";
|
||||
|
||||
private static final String LIST_GROUP_KEYS = "SELECT distinct GROUP_KEY as CREATED from %PREFIX%MESSAGE_GROUP where REGION=?";
|
||||
|
||||
@@ -409,6 +411,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
}
|
||||
});
|
||||
this.removeMessage(messageToRemove.getHeaders().getId());
|
||||
this.updateMessageGroup(groupKey);
|
||||
return getMessageGroup(groupId);
|
||||
}
|
||||
|
||||
@@ -458,12 +461,14 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
ps.setString(4, region);
|
||||
}
|
||||
});
|
||||
this.updateMessageGroup(groupKey);
|
||||
}
|
||||
|
||||
public Message<?> pollMessageFromGroup(final Object groupId) {
|
||||
String key = getKey(groupId);
|
||||
|
||||
return jdbcTemplate.query(getQuery(LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region },
|
||||
|
||||
Message<?> message = jdbcTemplate.query(getQuery(LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region },
|
||||
new ResultSetExtractor<Message<?>>() {
|
||||
public Message<?> extractData(ResultSet rs)
|
||||
throws SQLException, DataAccessException {
|
||||
@@ -480,25 +485,10 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
return null;
|
||||
}
|
||||
});
|
||||
this.updateMessageGroup(key);
|
||||
return message;
|
||||
}
|
||||
|
||||
private List<UUID> getMessageIdsForGroup(Object groupId){
|
||||
String key = getKey(groupId);
|
||||
|
||||
final List<UUID> messageIds = new ArrayList<UUID>();
|
||||
|
||||
jdbcTemplate.query(getQuery(LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region },
|
||||
new RowCallbackHandler() {
|
||||
|
||||
public void processRow(ResultSet rs) throws SQLException {
|
||||
messageIds.add(UUID.fromString(rs.getString(1)));
|
||||
}
|
||||
|
||||
}
|
||||
);
|
||||
return messageIds;
|
||||
}
|
||||
|
||||
|
||||
public Iterator<MessageGroup> iterator() {
|
||||
|
||||
final Iterator<String> iterator = jdbcTemplate.query(getQuery(LIST_GROUP_KEYS), new Object[] { region },
|
||||
@@ -521,6 +511,36 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
};
|
||||
|
||||
}
|
||||
|
||||
private void updateMessageGroup(final String groupId){
|
||||
jdbcTemplate.update(getQuery(UPDATE_GROUP), new PreparedStatementSetter() {
|
||||
public void setValues(PreparedStatement ps) throws SQLException {
|
||||
logger.debug("Updating MessageGroup: " + groupId);
|
||||
ps.setTimestamp(1, new Timestamp(System.currentTimeMillis()));
|
||||
ps.setString(2, groupId);
|
||||
ps.setString(3, region);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private List<UUID> getMessageIdsForGroup(Object groupId){
|
||||
String key = getKey(groupId);
|
||||
|
||||
final List<UUID> messageIds = new ArrayList<UUID>();
|
||||
|
||||
jdbcTemplate.query(getQuery(LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region },
|
||||
new RowCallbackHandler() {
|
||||
|
||||
public void processRow(ResultSet rs) throws SQLException {
|
||||
messageIds.add(UUID.fromString(rs.getString(1)));
|
||||
}
|
||||
|
||||
}
|
||||
);
|
||||
return messageIds;
|
||||
}
|
||||
|
||||
|
||||
|
||||
private String getKey(Object input) {
|
||||
return input == null ? null : UUIDConverter.getUUID(input).toString();
|
||||
|
||||
@@ -196,6 +196,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
Assert.notNull(groupId, "'groupId' must not be null");
|
||||
Assert.notNull(messageToRemove, "'messageToRemove' must not be null");
|
||||
this.removeMessage(messageToRemove.getHeaders().getId());
|
||||
this.updateGroup(groupId);
|
||||
return this.getMessageGroup(groupId);
|
||||
}
|
||||
|
||||
@@ -222,12 +223,14 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
Update update = Update.update(GROUP_COMPLETE_KEY, true);
|
||||
Query q = whereGroupIdIs(groupId);
|
||||
this.template.updateFirst(q, update, this.collectionName);
|
||||
this.updateGroup(groupId);
|
||||
}
|
||||
|
||||
public void setLastReleasedSequenceNumberForGroup(Object groupId, int sequenceNumber) {
|
||||
Update update = Update.update(LAST_RELEASED_SEQUENCE_NUMBER, sequenceNumber);
|
||||
Query q = whereGroupIdIs(groupId);
|
||||
this.template.updateFirst(q, update, this.collectionName);
|
||||
this.updateGroup(groupId);
|
||||
}
|
||||
|
||||
public Message<?> pollMessageFromGroup(Object groupId) {
|
||||
@@ -239,7 +242,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
message = messageWrappers.get(0).getMessage();
|
||||
this.removeMessageFromGroup(groupId, message);
|
||||
}
|
||||
|
||||
this.updateGroup(groupId);
|
||||
return message;
|
||||
}
|
||||
|
||||
@@ -266,6 +269,12 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
q.sort().on(CREATED_DATE, Order.ASCENDING);
|
||||
return q;
|
||||
}
|
||||
|
||||
private void updateGroup(Object groupId) {
|
||||
Update update = Update.update(GROUP_UPDATE_TIMESTAMP_KEY, System.currentTimeMillis());
|
||||
Query q = whereGroupIdIs(groupId);
|
||||
this.template.updateFirst(q, update, this.collectionName);
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user