INT-3037 Fix JDBC MS Discard After Completion

INT-3037 - commented out sizing component

unit test to support

cleanup

updated commit as per Artem's comments May 29
This commit is contained in:
willschipp
2013-05-28 14:58:37 -04:00
committed by Gary Russell
parent 077f2b24ea
commit 6140aa7b2c
2 changed files with 83 additions and 9 deletions

View File

@@ -75,6 +75,7 @@ import org.springframework.util.StringUtils;
* @author Oleg Zhurakousky
* @author Matt Stine
* @author Gunnar Hillert
* @author Will Schipp
*
* @since 2.0
*/
@@ -293,6 +294,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
Assert.state(jdbcTemplate != null, "A DataSource or JdbcTemplate must be provided");
}
@Override
public Message<?> removeMessage(UUID id) {
Message<?> message = getMessage(id);
if (message == null) {
@@ -306,11 +308,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return null;
}
@Override
@ManagedAttribute
public long getMessageCount() {
return jdbcTemplate.queryForInt(getQuery(Query.GET_MESSAGE_COUNT), region);
}
@Override
public Message<?> getMessage(UUID id) {
List<Message<?>> list = jdbcTemplate.query(getQuery(Query.GET_MESSAGE), new Object[] { getKey(id), region }, mapper);
if (list.isEmpty()) {
@@ -319,6 +323,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return list.get(0);
}
@Override
@SuppressWarnings({ "rawtypes", "unchecked" })
public <T> Message<T> addMessage(final Message<T> message) {
if (message.getHeaders().containsKey(SAVED_KEY)) {
@@ -342,6 +347,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
final byte[] messageBytes = serializer.convert(result);
jdbcTemplate.update(getQuery(Query.CREATE_MESSAGE), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Inserting message with id key=" + messageId);
@@ -355,6 +361,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return result;
}
@Override
public MessageGroup addMessageToGroup(Object groupId, Message<?> message) {
final String groupKey = getKey(groupId);
final String messageId = getKey(message.getHeaders().getId());
@@ -381,6 +388,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
this.addMessage(message);
jdbcTemplate.update(getQuery(Query.CREATE_GROUP_TO_MESSAGE), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Inserting message with id key=" + messageId + " and created date=" + createdDate);
@@ -406,12 +414,14 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return jdbcTemplate.queryForInt(getQuery(Query.COUNT_ALL_MESSAGES_IN_GROUPS), region);
}
@Override
@ManagedAttribute
public int messageGroupSize(Object groupId) {
String key = getKey(groupId);
return jdbcTemplate.queryForInt(getQuery(Query.COUNT_ALL_MESSAGES_IN_GROUP), key, region);
}
@Override
public MessageGroup getMessageGroup(Object groupId) {
String key = getKey(groupId);
final AtomicReference<Date> createDate = new AtomicReference<Date>();
@@ -421,12 +431,9 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
List<Message<?>> messages = jdbcTemplate.query(getQuery(Query.LIST_MESSAGES_BY_GROUP_KEY), new Object[] { key, region }, mapper);
if (messages.size() == 0){
return new SimpleMessageGroup(groupId);
}
jdbcTemplate.query(getQuery(Query.GET_GROUP_INFO), new Object[] { key, region},
new RowCallbackHandler() {
@Override
public void processRow(ResultSet rs) throws SQLException {
updateDate.set(rs.getTimestamp("UPDATED_DATE"));
@@ -459,11 +466,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return messageGroup;
}
@Override
public MessageGroup removeMessageFromGroup(Object groupId, Message<?> messageToRemove) {
final String groupKey = getKey(groupId);
final String messageId = getKey(messageToRemove.getHeaders().getId());
jdbcTemplate.update(getQuery(Query.REMOVE_MESSAGE_FROM_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Removing message from group with group key=" + groupKey);
@@ -478,6 +487,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return getMessageGroup(groupId);
}
@Override
public void removeMessageGroup(Object groupId) {
final String groupKey = getKey(groupId);
@@ -487,6 +497,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
}
jdbcTemplate.update(getQuery(Query.REMOVE_GROUP_TO_MESSAGE_JOIN), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Removing relationships for the group with group key=" + groupKey);
@@ -497,6 +508,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
});
jdbcTemplate.update(getQuery(Query.DELETE_MESSAGE_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Marking messages with group key=" + groupKey);
@@ -507,11 +519,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
});
}
@Override
public void completeGroup(Object groupId) {
final long updatedDate = System.currentTimeMillis();
final String groupKey = getKey(groupId);
jdbcTemplate.update(getQuery(Query.COMPLETE_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Completing MessageGroup: " + groupKey);
@@ -523,12 +537,14 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
});
}
@Override
public void setLastReleasedSequenceNumberForGroup(Object groupId, final int sequenceNumber) {
Assert.notNull(groupId, "'groupId' must not be null");
final long updatedDate = System.currentTimeMillis();
final String groupKey = getKey(groupId);
jdbcTemplate.update(getQuery(Query.UPDATE_LAST_RELEASED_SEQUENCE), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Updating the sequence number of the last released Message in the MessageGroup: " + groupKey);
@@ -542,6 +558,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
this.updateMessageGroup(groupKey);
}
@Override
public Message<?> pollMessageFromGroup(Object groupId) {
String key = getKey(groupId);
@@ -552,6 +569,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return polledMessage;
}
@Override
public Iterator<MessageGroup> iterator() {
final Iterator<String> iterator = jdbcTemplate.query(getQuery(Query.LIST_GROUP_KEYS), new Object[] { region },
@@ -559,14 +577,17 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
return new Iterator<MessageGroup>() {
@Override
public boolean hasNext() {
return iterator.hasNext();
}
@Override
public MessageGroup next() {
return getMessageGroup(iterator.next());
}
@Override
public void remove() {
throw new UnsupportedOperationException("Cannot remove MessageGroup from this iterator.");
}
@@ -620,6 +641,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
private void doCreateMessageGroup(final String groupKey, final Timestamp createdDate){
jdbcTemplate.update(getQuery(Query.CREATE_MESSAGE_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Creating message group with id key=" + groupKey + " and created date=" + createdDate);
@@ -634,6 +656,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
private void doUpdateMessageGroup(final String groupKey, final Timestamp updatedDate){
jdbcTemplate.update(getQuery(Query.UPDATE_MESSAGE_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Updating message group with id key=" + groupKey + " and updated date=" + updatedDate);
@@ -647,6 +670,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
private void updateMessageGroup(final String groupId){
jdbcTemplate.update(getQuery(Query.UPDATE_GROUP), new PreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps) throws SQLException {
if (logger.isDebugEnabled()){
logger.debug("Updating MessageGroup: " + groupId);
@@ -666,6 +690,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
jdbcTemplate.query(getQuery(Query.LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region },
new RowCallbackHandler() {
@Override
public void processRow(ResultSet rs) throws SQLException {
messageIds.add(UUID.fromString(rs.getString(1)));
}
@@ -686,6 +711,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
*/
private class MessageMapper implements RowMapper<Message<?>> {
@Override
public Message<?> mapRow(ResultSet rs, int rowNum) throws SQLException {
Message<?> message = (Message<?>) deserializer.convert(lobHandler.getBlobAsBytes(rs, "MESSAGE_BYTES"));
return message;

View File

@@ -65,6 +65,7 @@ import org.springframework.transaction.annotation.Transactional;
* @author Gunnar Hillert
* @author Artem Bilan
* @author Gary Russell
* @author Will Schipp
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
@@ -137,12 +138,14 @@ public class JdbcMessageStoreTests {
public void testSerializer() throws Exception {
// N.B. these serializers are not realistic (just for test purposes)
messageStore.setSerializer(new Serializer<Message<?>>() {
@Override
public void serialize(Message<?> object, OutputStream outputStream) throws IOException {
outputStream.write(((Message<?>) object).getPayload().toString().getBytes());
outputStream.flush();
}
});
messageStore.setDeserializer(new Deserializer<GenericMessage<String>>() {
@Override
public GenericMessage<String> deserialize(InputStream inputStream) throws IOException {
BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream));
return new GenericMessage<String>(reader.readLine());
@@ -320,6 +323,7 @@ public class JdbcMessageStoreTests {
Message<String> message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build();
messageStore.addMessageToGroup(groupId, message);
messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() {
@Override
public void execute(MessageGroupStore messageGroupStore, MessageGroup group) {
messageGroupStore.removeMessageGroup(group.getGroupId());
}
@@ -343,6 +347,7 @@ public class JdbcMessageStoreTests {
messageStore.setTimeoutOnIdle(true);
messageStore.addMessageToGroup(groupId, message);
messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() {
@Override
public void execute(MessageGroupStore messageGroupStore, MessageGroup group) {
messageGroupStore.removeMessageGroup(group.getGroupId());
}
@@ -424,8 +429,8 @@ public class JdbcMessageStoreTests {
LOG.info("messageFromGroup1: " + messageFromGroup1.getHeaders().getId() + "; Sequence #: " + messageFromGroup1.getHeaders().getSequenceNumber());
LOG.info("messageFromGroup2: " + messageFromGroup2.getHeaders().getId() + "; Sequence #: " + messageFromGroup2.getHeaders().getSequenceNumber());
assertEquals(Integer.valueOf(1), (Integer) messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(2), (Integer) messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(1), messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(2), messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
}
@@ -466,8 +471,51 @@ public class JdbcMessageStoreTests {
LOG.info("messageFromRegion1: " + messageFromRegion1.getHeaders().getId() + "; Sequence #: " + messageFromRegion1.getHeaders().getSequenceNumber());
LOG.info("messageFromRegion2: " + messageFromRegion2.getHeaders().getId() + "; Sequence #: " + messageFromRegion2.getHeaders().getSequenceNumber());
assertEquals(Integer.valueOf(1), (Integer) messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(2), (Integer) messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(1), messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
assertEquals(Integer.valueOf(2), messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER));
}
}
@Test
@Transactional
public void testCompletedNotExpiredGroupINT3037() throws Exception {
/*
* based on the aggregator scenario as follows;
*
* send three messages in
* 1 of 2
* 2 of 2
* 2 of 2 (last again)
*
* expected behavior is that the LAST message (2 of 2 repeat) should be on the discard channel
* (discard behavior performed by the AbstractCorrelatingMessageHandler.handleMessageInternal)
*/
final JdbcMessageStore messageStore = new JdbcMessageStore(dataSource);
//init
String groupId = "group";
//build the messages
Message<?> oneOfTwo = MessageBuilder.withPayload("hello").setSequenceNumber(1).setSequenceSize(2).setCorrelationId(groupId).build();
Message<?> twoOfTwo = MessageBuilder.withPayload("world").setSequenceNumber(2).setSequenceSize(2).setCorrelationId(groupId).build();
//add to the messageStore
messageStore.addMessageToGroup(groupId, oneOfTwo);
messageStore.addMessageToGroup(groupId, twoOfTwo);
//check that 2 messages are there
assertTrue(messageStore.getMessageGroupCount() == 1);
assertTrue(messageStore.getMessageCount() == 2);
//retrieve the group (like in the aggregator)
MessageGroup messageGroup = messageStore.getMessageGroup(groupId);
//'complete' the group
messageStore.completeGroup(messageGroup.getGroupId());
//now clear the messages
for (Message<?> message : messageGroup.getMessages()) {
messageStore.removeMessageFromGroup(groupId, message);
}//end for
//'add' the other message --> emulated by getting the messageGroup
messageGroup = messageStore.getMessageGroup(groupId);
//should be marked 'complete' --> old behavior it would not
assertTrue(messageGroup.isComplete());
}
}