Some code clean up for ChMessStoreQueryProvider

* Add a version for deprecation on `AbstractChannelMessageStoreQueryProvider`
* Remove dead links from `ChannelMessageStoreQueryProvider` impl JavaDocs
* Use text blocks for long query strings
* Move `SELECT_COMMON` constant to the `ChannelMessageStoreQueryProvider` interface
This commit is contained in:
Artem Bilan
2023-09-01 11:15:09 -04:00
parent 0a95091331
commit eef070a3d9
9 changed files with 113 additions and 107 deletions

View File

@@ -23,8 +23,8 @@ package org.springframework.integration.jdbc.store.channel;
* @author Adama Sorho
*
* @since 2.2
* @deprecated in favor of default methods in ChannelMessageStoreQueryProvider
* @deprecated since 6.2 in favor of default methods in {@link ChannelMessageStoreQueryProvider}
*/
@Deprecated
@Deprecated(since = "6.2", forRemoval = true)
public abstract class AbstractChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
}

View File

@@ -30,60 +30,37 @@ package org.springframework.integration.jdbc.store.channel;
*/
public interface ChannelMessageStoreQueryProvider {
String SELECT_COMMON = """
SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES
from %PREFIX%CHANNEL_MESSAGE
where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region\s
""";
/**
* Get the query used to retrieve a count of all messages currently persisted
* for a channel.
*
* @return Sql Query
* @return query string
*/
default String getCountAllMessagesInGroupQuery() {
return "SELECT COUNT(MESSAGE_ID) from %PREFIX%CHANNEL_MESSAGE where GROUP_KEY=? and REGION=?";
}
/**
* Get the query used to retrieve the oldest message for a channel excluding
* messages that match the provided message ids.
*
* @return Sql Query
*/
String getPollFromGroupExcludeIdsQuery();
/**
* Get the query used to retrieve the oldest message for a channel.
*
* @return Sql Query
*/
String getPollFromGroupQuery();
/**
* Get the query used to retrieve the oldest message by priority for a channel excluding
* messages that match the provided message ids.
*
* @return Sql Query
*/
String getPriorityPollFromGroupExcludeIdsQuery();
/**
* Get the query used to retrieve the oldest message by priority for a channel.
*
* @return Sql Query
*/
String getPriorityPollFromGroupQuery();
/**
* Query that retrieves a message for the provided message id, channel and
* region.
*
* @return Sql Query
* @return query string
*/
default String getMessageQuery() {
return "SELECT MESSAGE_ID, CREATED_DATE, MESSAGE_BYTES from %PREFIX%CHANNEL_MESSAGE where MESSAGE_ID=? and GROUP_KEY=? and REGION=?";
return """
SELECT MESSAGE_ID, CREATED_DATE, MESSAGE_BYTES
from %PREFIX%CHANNEL_MESSAGE
where MESSAGE_ID=? and GROUP_KEY=? and REGION=?
""";
}
/**
* Query that retrieve a count of all messages for a region.
*
* @return Sql Query
* @return query string
*/
default String getMessageCountForRegionQuery() {
return "SELECT COUNT(MESSAGE_ID) from %PREFIX%CHANNEL_MESSAGE where REGION=?";
@@ -91,8 +68,7 @@ public interface ChannelMessageStoreQueryProvider {
/**
* Query to delete a single message from the database.
*
* @return Sql Query
* @return query string
*/
default String getDeleteMessageQuery() {
return "DELETE from %PREFIX%CHANNEL_MESSAGE where MESSAGE_ID=? and GROUP_KEY=? and REGION=?";
@@ -100,21 +76,53 @@ public interface ChannelMessageStoreQueryProvider {
/**
* Query to add a single message to the database.
*
* @return Sql Query
* @return query string
*/
default String getCreateMessageQuery() {
return "INSERT into %PREFIX%CHANNEL_MESSAGE(MESSAGE_ID, GROUP_KEY, REGION, CREATED_DATE, MESSAGE_PRIORITY, MESSAGE_BYTES)"
+ " values (?, ?, ?, ?, ?, ?)";
return """
INSERT into %PREFIX%CHANNEL_MESSAGE(
MESSAGE_ID,
GROUP_KEY,
REGION,
CREATED_DATE,
MESSAGE_PRIORITY,
MESSAGE_BYTES)
values (?, ?, ?, ?, ?, ?)
""";
}
/**
* Query to delete all messages that belong to a specific channel.
*
* @return Sql Query
* @return query string
*/
default String getDeleteMessageGroupQuery() {
return "DELETE from %PREFIX%CHANNEL_MESSAGE where GROUP_KEY=? and REGION=?";
}
/**
* Get the query used to retrieve the oldest message for a channel excluding
* messages that match the provided message ids.
* @return query string
*/
String getPollFromGroupExcludeIdsQuery();
/**
* Get the query used to retrieve the oldest message for a channel.
* @return query string
*/
String getPollFromGroupQuery();
/**
* Get the query used to retrieve the oldest message by priority for a channel excluding
* messages that match the provided message ids.
* @return query string
*/
String getPriorityPollFromGroupExcludeIdsQuery();
/**
* Get the query used to retrieve the oldest message by priority for a channel.
* @return query string
*/
String getPriorityPollFromGroupQuery();
}

View File

@@ -21,17 +21,11 @@ package org.springframework.integration.jdbc.store.channel;
* @author Artem Bilan
* @author Gary Russell
* @author Adama Sorho
* @since 2.2
*
* https://blogs.oracle.com/kah/entry/derby_10_5_preview_fetch
* @since 2.2
*/
public class DerbyChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON

View File

@@ -22,16 +22,12 @@ package org.springframework.integration.jdbc.store.channel;
* @author Manuel Jordan
* @author Gary Russell
* @author Adama Sorho
*
* @since 4.3
*
*/
public class H2ChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getCreateMessageQuery() {
return "INSERT into %PREFIX%CHANNEL_MESSAGE(MESSAGE_ID, GROUP_KEY, REGION, CREATED_DATE, MESSAGE_PRIORITY, " +

View File

@@ -20,26 +20,32 @@ package org.springframework.integration.jdbc.store.channel;
* @author Gunnar Hillert
* @author Artem Bilan
* @author Adama Sorho
*
* @since 2.2
*
*/
public class HsqlChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getCreateMessageQuery() {
return "INSERT into %PREFIX%CHANNEL_MESSAGE(MESSAGE_ID, GROUP_KEY, REGION, CREATED_DATE, MESSAGE_PRIORITY, MESSAGE_SEQUENCE, MESSAGE_BYTES)"
+ " values (?, ?, ?, ?, ?, NEXT VALUE FOR %PREFIX%MESSAGE_SEQ, ?)";
return """
INSERT into %PREFIX%CHANNEL_MESSAGE(
MESSAGE_ID,
GROUP_KEY,
REGION,
CREATED_DATE,
MESSAGE_PRIORITY,
MESSAGE_SEQUENCE,
MESSAGE_BYTES)
values (?, ?, ?, ?, ?, NEXT VALUE FOR %PREFIX%MESSAGE_SEQ, ?)
""";
}
@Override
public String getPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON +
"and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) order by CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1";
"and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " +
"order by CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1";
}
@Override

View File

@@ -25,16 +25,12 @@ package org.springframework.integration.jdbc.store.channel;
*/
public class MySqlChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON
+ "and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) "
+ "order by CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED";
+ "order by CREATED_DATE, MESSAGE_SEQUENCE " +
"LIMIT 1 FOR UPDATE SKIP LOCKED";
}
@Override
@@ -47,13 +43,15 @@ public class MySqlChannelMessageStoreQueryProvider implements ChannelMessageStor
public String getPriorityPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON +
"and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " +
"order by MESSAGE_PRIORITY DESC, CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED";
"order by MESSAGE_PRIORITY DESC, CREATED_DATE, MESSAGE_SEQUENCE " +
"LIMIT 1 FOR UPDATE SKIP LOCKED";
}
@Override
public String getPriorityPollFromGroupQuery() {
return SELECT_COMMON +
"order by MESSAGE_PRIORITY DESC, CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED";
"order by MESSAGE_PRIORITY DESC, CREATED_DATE, MESSAGE_SEQUENCE " +
"LIMIT 1 FOR UPDATE SKIP LOCKED";
}
}

View File

@@ -20,10 +20,8 @@ package org.springframework.integration.jdbc.store.channel;
* Contains Oracle-specific queries for the
* {@link org.springframework.integration.jdbc.store.JdbcChannelMessageStore}. Please
* ensure that the used {@link org.springframework.jdbc.core.JdbcTemplate}'s fetchSize
* property is <code>1</code>.
* property is {@code 1}.
* <p>
* Fore more details, please see:
* https://stackoverflow.com/questions/6117254/force-oracle-to-return-top-n-rows-with-skip-locked
*
* @author Gunnar Hillert
* @author Artem Bilan
@@ -33,11 +31,6 @@ package org.springframework.integration.jdbc.store.channel;
*/
public class OracleChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getCreateMessageQuery() {
return "INSERT into %PREFIX%CHANNEL_MESSAGE(MESSAGE_ID, GROUP_KEY, REGION, CREATED_DATE, MESSAGE_PRIORITY, "
@@ -60,19 +53,27 @@ public class OracleChannelMessageStoreQueryProvider implements ChannelMessageSto
@Override
public String getPriorityPollFromGroupExcludeIdsQuery() {
return "SELECT /*+ INDEX(%PREFIX%CHANNEL_MESSAGE %PREFIX%CHANNEL_MSG_PRIORITY_IDX) */ " +
"%PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES from %PREFIX%CHANNEL_MESSAGE " +
"where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region " +
"and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " +
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE FOR UPDATE SKIP LOCKED";
return """
SELECT /*+ INDEX(%PREFIX%CHANNEL_MESSAGE %PREFIX%CHANNEL_MSG_PRIORITY_IDX) */
%PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES
from %PREFIX%CHANNEL_MESSAGE
where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key
and %PREFIX%CHANNEL_MESSAGE.REGION = :region
and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids)
order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE FOR UPDATE SKIP LOCKED
""";
}
@Override
public String getPriorityPollFromGroupQuery() {
return "SELECT /*+ INDEX(%PREFIX%CHANNEL_MESSAGE %PREFIX%CHANNEL_MSG_PRIORITY_IDX) */ " +
"%PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES from %PREFIX%CHANNEL_MESSAGE " +
"where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region " +
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE FOR UPDATE SKIP LOCKED";
return """
SELECT /*+ INDEX(%PREFIX%CHANNEL_MESSAGE %PREFIX%CHANNEL_MSG_PRIORITY_IDX) */
%PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES
from %PREFIX%CHANNEL_MESSAGE
where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key
and %PREFIX%CHANNEL_MESSAGE.REGION = :region
order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE FOR UPDATE SKIP LOCKED
""";
}
}

View File

@@ -25,11 +25,6 @@ package org.springframework.integration.jdbc.store.channel;
*/
public class PostgresChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON
@@ -47,13 +42,15 @@ public class PostgresChannelMessageStoreQueryProvider implements ChannelMessageS
public String getPriorityPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON +
"and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " +
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED";
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE " +
"LIMIT 1 FOR UPDATE SKIP LOCKED";
}
@Override
public String getPriorityPollFromGroupQuery() {
return SELECT_COMMON +
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED";
"order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE " +
"LIMIT 1 FOR UPDATE SKIP LOCKED";
}
}

View File

@@ -18,17 +18,15 @@ package org.springframework.integration.jdbc.store.channel;
/**
* Channel message store query provider for Microsoft SQL Server / Azure SQL database.
*
* @author Sundara Balaji
* @author Adama Sorho
* @author Artem Bilan
*
* @since 5.1
*/
public class SqlServerChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
private static final String SELECT_COMMON =
"SELECT TOP 1 %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID, %PREFIX%CHANNEL_MESSAGE.MESSAGE_BYTES "
+ "from %PREFIX%CHANNEL_MESSAGE "
+ "where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key and %PREFIX%CHANNEL_MESSAGE.REGION = :region ";
@Override
public String getPollFromGroupExcludeIdsQuery() {
return SELECT_COMMON +
@@ -56,9 +54,17 @@ public class SqlServerChannelMessageStoreQueryProvider implements ChannelMessage
@Override
public String getCreateMessageQuery() {
return "INSERT into %PREFIX%CHANNEL_MESSAGE(MESSAGE_ID, GROUP_KEY, REGION, CREATED_DATE, MESSAGE_PRIORITY, "
+ "MESSAGE_SEQUENCE, MESSAGE_BYTES)"
+ " values (?, ?, ?, ?, ?,(NEXT VALUE FOR %PREFIX%MESSAGE_SEQ), ?)";
return """
INSERT into %PREFIX%CHANNEL_MESSAGE(
MESSAGE_ID,
GROUP_KEY,
REGION,
CREATED_DATE,
MESSAGE_PRIORITY,
MESSAGE_SEQUENCE,
MESSAGE_BYTES)
values (?, ?, ?, ?, ?,(NEXT VALUE FOR %PREFIX%MESSAGE_SEQ), ?)
""";
}
}