From bd9d591dbfe38dc60c4041b63306621357a1f83a Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 22 Nov 2011 04:16:17 -0500 Subject: [PATCH] INT-2253 added support for updating MessageGroup with lastModified timestamp on every touch --- .../store/AbstractKeyValueMessageStore.java | 6 +- .../store/MessageGroupMetadata.java | 4 ++ .../integration/jdbc/JdbcMessageStore.java | 58 +++++++++++++------ .../mongodb/store/MongoDbMessageStore.java | 11 +++- 4 files changed, 58 insertions(+), 21 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java index a41a94931f..803ee4e96a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java @@ -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); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java index d0dd87ff4b..a4d01b2cb6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java @@ -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; diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java index d67cfa185b..551e7c5b14 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java @@ -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>() { 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 getMessageIdsForGroup(Object groupId){ - String key = getKey(groupId); - - final List messageIds = new ArrayList(); - - 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 iterator() { final Iterator 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 getMessageIdsForGroup(Object groupId){ + String key = getKey(groupId); + + final List messageIds = new ArrayList(); + + 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(); diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java index 68d8cd7f84..bf4d26defa 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java @@ -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); + } /**