From 7c701ca5e634187894af157b219f344d0350fbc1 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 21 Nov 2017 12:19:55 -0500 Subject: [PATCH] INT-2576: Optimize JdbcMS.removeMessageGroup JIRA: https://jira.spring.io/browse/INT-2576 * Do not fetch message ids and then iterate over them to remove one by one - just use an appropriate `DELETE` query for all messages --- .../jdbc/store/JdbcMessageStore.java | 35 +++++++------------ 1 file changed, 12 insertions(+), 23 deletions(-) diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java index d34e2ef76b..298abfd45d 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java @@ -19,7 +19,6 @@ package org.springframework.integration.jdbc.store; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Timestamp; -import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Date; @@ -48,7 +47,6 @@ import org.springframework.integration.support.converter.WhiteListDeserializingC import org.springframework.integration.util.UUIDConverter; import org.springframework.jdbc.core.JdbcOperations; import org.springframework.jdbc.core.JdbcTemplate; -import org.springframework.jdbc.core.RowCallbackHandler; import org.springframework.jdbc.core.RowMapper; import org.springframework.jdbc.core.SingleColumnRowMapper; import org.springframework.jdbc.support.lob.DefaultLobHandler; @@ -105,12 +103,9 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa COUNT_ALL_MESSAGES_IN_GROUP("SELECT COUNT(MESSAGE_ID) from %PREFIX%GROUP_TO_MESSAGE where GROUP_KEY=? and REGION=?"), - LIST_MESSAGEIDS_BY_GROUP_KEY("select MESSAGE_ID, CREATED_DATE " + - "from %PREFIX%MESSAGE where MESSAGE_ID in (select MESSAGE_ID from %PREFIX%GROUP_TO_MESSAGE where GROUP_KEY=? and REGION=?) " + - "ORDER BY CREATED_DATE"), - LIST_MESSAGES_BY_GROUP_KEY("SELECT MESSAGE_ID, MESSAGE_BYTES, CREATED_DATE " + - "from %PREFIX%MESSAGE where MESSAGE_ID in (SELECT MESSAGE_ID from %PREFIX%GROUP_TO_MESSAGE where GROUP_KEY = ?) and REGION=? " + + "from %PREFIX%MESSAGE where MESSAGE_ID in " + + "(SELECT MESSAGE_ID from %PREFIX%GROUP_TO_MESSAGE where GROUP_KEY = ? and REGION = ?) and REGION = ? " + "ORDER BY CREATED_DATE"), POLL_FROM_GROUP("SELECT %PREFIX%MESSAGE.MESSAGE_ID, %PREFIX%MESSAGE.MESSAGE_BYTES from %PREFIX%MESSAGE " + @@ -145,6 +140,9 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa UPDATE_LAST_RELEASED_SEQUENCE("UPDATE %PREFIX%MESSAGE_GROUP set UPDATED_DATE=?, LAST_RELEASED_SEQUENCE=? where GROUP_KEY=? and REGION=?"), + DELETE_MESSAGES_FROM_GROUP("DELETE from %PREFIX%MESSAGE where MESSAGE_ID in " + + "(SELECT MESSAGE_ID from %PREFIX%GROUP_TO_MESSAGE where GROUP_KEY = ? and REGION = ?) and REGION = ?"), + DELETE_MESSAGE_GROUP("DELETE from %PREFIX%MESSAGE_GROUP where GROUP_KEY=? and REGION=?"), CREATE_GROUP_TO_MESSAGE("INSERT into %PREFIX%GROUP_TO_MESSAGE" + @@ -478,12 +476,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa @Override public void removeMessageGroup(Object groupId) { + String groupKey = getKey(groupId); - final String groupKey = getKey(groupId); - - for (UUID messageIds : this.getMessageIdsForGroup(groupId)) { - this.removeMessage(messageIds); - } + this.jdbcTemplate.update(getQuery(Query.DELETE_MESSAGES_FROM_GROUP), ps -> { + ps.setString(1, groupKey); + ps.setString(2, JdbcMessageStore.this.region); + ps.setString(3, JdbcMessageStore.this.region); + }); this.jdbcTemplate.update(getQuery(Query.REMOVE_GROUP_TO_MESSAGE_JOIN), ps -> { if (logger.isDebugEnabled()) { @@ -555,7 +554,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa @Override public Collection> getMessagesForGroup(Object groupId) { return this.jdbcTemplate.query(getQuery(Query.LIST_MESSAGES_BY_GROUP_KEY), this.mapper, getKey(groupId), - this.region); + this.region, this.region); } @Override @@ -664,16 +663,6 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa }); } - private List getMessageIdsForGroup(Object groupId) { - String key = getKey(groupId); - - final List messageIds = new ArrayList(); - - this.jdbcTemplate.query(getQuery(Query.LIST_MESSAGEIDS_BY_GROUP_KEY), - (RowCallbackHandler) rs -> messageIds.add(UUID.fromString(rs.getString(1))), key, this.region); - return messageIds; - } - private String getKey(Object input) { return input == null ? null : UUIDConverter.getUUID(input).toString(); }