|
20 | 20 | * @author Gunnar Hillert
|
21 | 21 | * @author Artem Bilan
|
22 | 22 | * @author Adama Sorho
|
| 23 | + * @author Johannes Edmeier |
23 | 24 | *
|
24 | 25 | * @since 2.2
|
25 | 26 | */
|
26 | 27 | public class PostgresChannelMessageStoreQueryProvider implements ChannelMessageStoreQueryProvider {
|
27 | 28 |
|
28 | 29 | @Override
|
29 | 30 | public String getPollFromGroupExcludeIdsQuery() {
|
30 |
| - return SELECT_COMMON |
31 |
| - + "and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " |
32 |
| - + "order by CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED"; |
| 31 | + return """ |
| 32 | + delete |
| 33 | + from %PREFIX%CHANNEL_MESSAGE |
| 34 | + where CTID = (select CTID |
| 35 | + from %PREFIX%CHANNEL_MESSAGE |
| 36 | + where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key |
| 37 | + and %PREFIX%CHANNEL_MESSAGE.REGION = :region |
| 38 | + and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) |
| 39 | + order by CREATED_DATE, MESSAGE_SEQUENCE |
| 40 | + limit 1 for update skip locked) |
| 41 | + returning MESSAGE_ID, MESSAGE_BYTES; |
| 42 | + """; |
33 | 43 | }
|
34 | 44 |
|
35 | 45 | @Override
|
36 | 46 | public String getPollFromGroupQuery() {
|
37 |
| - return SELECT_COMMON + |
38 |
| - "order by CREATED_DATE, MESSAGE_SEQUENCE LIMIT 1 FOR UPDATE SKIP LOCKED"; |
| 47 | + return """ |
| 48 | + delete |
| 49 | + from %PREFIX%CHANNEL_MESSAGE |
| 50 | + where CTID = (select CTID |
| 51 | + from %PREFIX%CHANNEL_MESSAGE |
| 52 | + where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key |
| 53 | + and %PREFIX%CHANNEL_MESSAGE.REGION = :region |
| 54 | + order by CREATED_DATE, MESSAGE_SEQUENCE |
| 55 | + limit 1 for update skip locked) |
| 56 | + returning MESSAGE_ID, MESSAGE_BYTES; |
| 57 | + """; |
39 | 58 | }
|
40 | 59 |
|
41 | 60 | @Override
|
42 | 61 | public String getPriorityPollFromGroupExcludeIdsQuery() {
|
43 |
| - return SELECT_COMMON + |
44 |
| - "and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) " + |
45 |
| - "order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE " + |
46 |
| - "LIMIT 1 FOR UPDATE SKIP LOCKED"; |
| 62 | + return """ |
| 63 | + delete |
| 64 | + from %PREFIX%CHANNEL_MESSAGE |
| 65 | + where CTID = (select CTID |
| 66 | + from %PREFIX%CHANNEL_MESSAGE |
| 67 | + where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key |
| 68 | + and %PREFIX%CHANNEL_MESSAGE.REGION = :region |
| 69 | + and %PREFIX%CHANNEL_MESSAGE.MESSAGE_ID not in (:message_ids) |
| 70 | + order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE |
| 71 | + limit 1 for update skip locked) |
| 72 | + returning MESSAGE_ID, MESSAGE_BYTES; |
| 73 | + """; |
47 | 74 | }
|
48 | 75 |
|
49 | 76 | @Override
|
50 | 77 | public String getPriorityPollFromGroupQuery() {
|
51 |
| - return SELECT_COMMON + |
52 |
| - "order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE " + |
53 |
| - "LIMIT 1 FOR UPDATE SKIP LOCKED"; |
| 78 | + return """ |
| 79 | + delete |
| 80 | + from %PREFIX%CHANNEL_MESSAGE |
| 81 | + where CTID = (select CTID |
| 82 | + from %PREFIX%CHANNEL_MESSAGE |
| 83 | + where %PREFIX%CHANNEL_MESSAGE.GROUP_KEY = :group_key |
| 84 | + and %PREFIX%CHANNEL_MESSAGE.REGION = :region |
| 85 | + order by MESSAGE_PRIORITY DESC NULLS LAST, CREATED_DATE, MESSAGE_SEQUENCE |
| 86 | + limit 1 for update skip locked) |
| 87 | + returning MESSAGE_ID, MESSAGE_BYTES; |
| 88 | + """; |
| 89 | + } |
| 90 | + |
| 91 | + @Override |
| 92 | + public boolean isSingleStatementForPoll() { |
| 93 | + return true; |
54 | 94 | }
|
55 | 95 |
|
56 | 96 | }
|
0 commit comments