Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -719,11 +719,14 @@ public Message[] browse() {
public void doBrowse(final List<Message> browseList, final int max) {
try {
if (topicStore != null) {
// a JDBC ack moves the last acked id of the subscription, so expiring a browsed
// message would also ack the non-expired messages before it
final boolean expireFromStore = topicStore.getType() == StoreType.JDBC;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This breaks modularity. There shouldn't be a check for an underlying store type and a subsequent change in logic. If there is a problem with the JDBC ack'n with browsed messages, we should fix that in the JDBC layer.

final List<Message> toExpire = new ArrayList<Message>();
topicStore.recover(new MessageRecoveryListener() {
@Override
public boolean recoverMessage(Message message) throws Exception {
if (message.isExpired()) {
if (!expireFromStore && message.isExpired()) {
toExpire.add(message);
}
browseList.add(message);
Expand Down Expand Up @@ -759,6 +762,9 @@ public boolean isDuplicate(MessageId id) {
}
}
}
if (expireFromStore) {
expireFromStore(topicStore);
}
Message[] msgs = subscriptionRecoveryPolicy.browse(getActiveMQDestination());
if (msgs != null) {
for (int i = 0; i < msgs.length && browseList.size() < max; i++) {
Expand Down Expand Up @@ -924,61 +930,67 @@ public boolean isDuplicate(MessageId ref) {
}
};

// Expires messages through TopicMessageStore.recoverExpired(), for the stores that support it
private void expireFromStore(TopicMessageStore store) throws Exception {
// get the sub keys that should be checked for expired messages
final var subs = durableSubscribers.entrySet().stream()
.filter(entry -> isEligibleForExpiration(entry.getValue()))
.map(Entry::getKey).collect(Collectors.toSet());

if (subs.isEmpty()) {
LOG.debug("Skipping topic expiration check for {}, no eligible subscriptions to check", destination);
return;
}

// For each eligible subscription, return the messages in the store that are expired
// The same message refs are shared between subs if duplicated so this is efficient
var expired = store.recoverExpired(subs, getMaxExpirePageSize(),
expiryListener);

final ConnectionContext connectionContext = createConnectionContext();
// Go through any expired messages and remove for each sub
for (Entry<SubscriptionKey, List<Message>> entry : expired.entrySet()) {
DurableTopicSubscription sub = durableSubscribers.get(entry.getKey());
List<Message> expiredMessages = entry.getValue();

// If the sub still exists and there are expired messages then process
if (sub != null && !expiredMessages.isEmpty()) {
// There's a small race condition here if the sub comes online,
// but it's not a big deal as at worst there maybe be duplicate acks for
// the expired message but the store can handle it
if (isEligibleForExpiration(sub)) {
expiredMessages.forEach(message -> {
message.setRegionDestination(Topic.this);
try {
// AMQ-9721 - Remove message from the cursor if it exists after
// loading from the store. Store recoverExpired() does not inc
// the ref count so we don't need to decrement here, but if
// the cursor finds its own copy in memory it will dec that ref.
sub.removePending(message);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
messageExpired(connectionContext, sub, message);
});
}
}
}
}

private final AtomicBoolean expiryTaskInProgress = new AtomicBoolean(false);
private final Runnable expireMessagesWork = () -> {
try {
final TopicMessageStore store = Topic.this.topicStore;
if (store != null && store.getType() == StoreType.KAHADB) {
if (store.getMessageCount() == 0) {
if (store != null && (store.getType() == StoreType.KAHADB || store.getType() == StoreType.JDBC)) {
// the message count is a cheap index lookup for KahaDB but a COUNT query for JDBC
if (store.getType() == StoreType.KAHADB && store.getMessageCount() == 0) {
LOG.debug("Skipping topic expiration check for {}, store size is 0", destination);
return;
}

// get the sub keys that should be checked for expired messages
final var subs = durableSubscribers.entrySet().stream()
.filter(entry -> isEligibleForExpiration(entry.getValue()))
.map(Entry::getKey).collect(Collectors.toSet());

if (subs.isEmpty()) {
LOG.debug("Skipping topic expiration check for {}, no eligible subscriptions to check", destination);
return;
}

// For each eligible subscription, return the messages in the store that are expired
// The same message refs are shared between subs if duplicated so this is efficient
var expired = store.recoverExpired(subs, getMaxExpirePageSize(),
expiryListener);

final ConnectionContext connectionContext = createConnectionContext();
// Go through any expired messages and remove for each sub
for (Entry<SubscriptionKey, List<Message>> entry : expired.entrySet()) {
DurableTopicSubscription sub = durableSubscribers.get(entry.getKey());
List<Message> expiredMessages = entry.getValue();

// If the sub still exists and there are expired messages then process
if (sub != null && !expiredMessages.isEmpty()) {
// There's a small race condition here if the sub comes online,
// but it's not a big deal as at worst there maybe be duplicate acks for
// the expired message but the store can handle it
if (isEligibleForExpiration(sub)) {
expiredMessages.forEach(message -> {
message.setRegionDestination(Topic.this);
try {
// AMQ-9721 - Remove message from the cursor if it exists after
// loading from the store. Store recoverExpired() does not inc
// the ref count so we don't need to decrement here, but if
// the cursor finds its own copy in memory it will dec that ref.
sub.removePending(message);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
messageExpired(connectionContext, sub, message);
});
}
}
}
expireFromStore(store);
} else {
// If not KahaDB, fall back to the legacy browse method because
// If not KahaDB or JDBC, fall back to the legacy browse method because
// the recoverExpired() method is not supported
doBrowse(new InsertionCountList<>(), getMaxExpirePageSize());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,9 @@ void doRecoverSubscription(TransactionContext c, ActiveMQDestination destination
void doRecoverNextMessages(TransactionContext c, ActiveMQDestination destination, String clientId, String subscriptionName, long seq, long priority, int maxReturned,
JDBCMessageRecoveryListener listener) throws Exception;

void doRecoverExpired(TransactionContext c, ActiveMQDestination destination, String clientId, String subscriptionName, long now, int maxReturned,
boolean isPrioritizedMessages, JDBCMessageRecoveryListener listener) throws Exception;

void doRecoverNextMessagesWithPriority(TransactionContext c, ActiveMQDestination destination, String clientId, String subscriptionName, long seq, long priority, int maxReturned,
JDBCMessageRecoveryListener listener) throws Exception;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@

import java.io.IOException;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.LinkedHashMap;
Expand Down Expand Up @@ -362,10 +364,57 @@ public void addSubscription(SubscriptionInfo subscriptionInfo, boolean retroacti
}
}

/**
* Returns, for each subscription, the expired messages it has not acked yet, oldest first.
* The ack table only keeps the last acked id of a subscription, so only the run of expired
* messages right after it is returned: acking a later expired message would also ack the
* non-expired messages before it. At most {@code max} distinct messages are loaded.
*/
@Override
public Map<SubscriptionKey, List<Message>> recoverExpired(Set<SubscriptionKey> subs, int max,
MessageRecoveryListener listener) {
throw new UnsupportedOperationException("recoverExpired not supported");
MessageRecoveryListener listener) throws Exception {
final Map<SubscriptionKey, List<Message>> expired = new HashMap<>();
// a message pending for several subscriptions is unmarshalled and passed to the listener once
final Map<Long, Message> loaded = new HashMap<>();
final long now = System.currentTimeMillis();
TransactionContext c = persistenceAdapter.getTransactionContext();
try {
for (SubscriptionKey sub : subs) {
// no check on max here: messages already loaded for another subscription cost nothing
if (!listener.hasSpace()) {
break;
}
adapter.doRecoverExpired(c, destination, sub.getClientId(), sub.getSubscriptionName(), now, max,
isPrioritizedMessages(), new JDBCMessageRecoveryListener() {
@Override
public boolean recoverMessage(long sequenceId, byte[] data) throws Exception {
Message msg = loaded.get(sequenceId);
if (msg == null) {
if (loaded.size() >= max || !listener.hasSpace()) {
return false;
}
msg = (Message) wireFormat.unmarshal(new ByteSequence(data));
msg.getMessageId().setBrokerSequenceId(sequenceId);
loaded.put(sequenceId, msg);
listener.recoverMessage(msg);
}
expired.computeIfAbsent(sub, k -> new ArrayList<>()).add(msg);
return true;
}

@Override
public boolean recoverMessageReference(String reference) {
return false;
}
});
}
} catch (SQLException e) {
JDBCPersistenceAdapter.log("JDBC Failure: ", e);
throw IOExceptionSupport.create("Failed to recover expired messages for: " + destination + ". Reason: " + e, e);
} finally {
c.close();
}
return expired;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,10 @@ public class Statements {
private String findAllDurableSubMessagesStatement;
private String findDurableSubMessagesStatement;
private String findDurableSubMessagesByPriorityStatement;
private String findExpiredDurableSubRangeStatement;
private String findExpiredDurableSubRangeByPriorityStatement;
private String findExpiredMessagesInRangeStatement;
private String findExpiredMessagesInRangeByPriorityStatement;
private String findAllDestinationsStatement;
private String removeAllMessagesStatement;
private String removeAllSubscriptionsStatement;
Expand Down Expand Up @@ -336,6 +340,74 @@ public String getFindDurableSubMessagesByPriorityStatement() {
return findDurableSubMessagesByPriorityStatement;
}

/**
* For each ack row of a durable subscription: the last acked id, the priority, the id of the
* first pending message that is not expired or is still part of a transaction (null if there
* is none) and the highest message id of the destination. The ack table only records the last
* acked id, so acking a message implicitly acks every earlier one: only the expired messages
* between the last acked id and that first blocking message can be expired.
* Parameters: now, container, client id, sub name.
*/
public String getFindExpiredDurableSubRangeStatement() {
if (findExpiredDurableSubRangeStatement == null) {
findExpiredDurableSubRangeStatement = "SELECT D.LAST_ACKED_ID, D.PRIORITY,"
+ " (SELECT MIN(B.ID) FROM " + getFullMessageTableName() + " B"
+ " WHERE B.CONTAINER=D.CONTAINER AND B.ID > D.LAST_ACKED_ID"
+ " AND (B.XID IS NOT NULL OR B.EXPIRATION <= 0 OR B.EXPIRATION >= ?)),"
+ " (SELECT MAX(B.ID) FROM " + getFullMessageTableName() + " B WHERE B.CONTAINER=D.CONTAINER)"
+ " FROM " + getFullAckTableName() + " D"
+ " WHERE D.CONTAINER=? AND D.CLIENT_ID=? AND D.SUB_NAME=? AND D.XID IS NULL";
}
return findExpiredDurableSubRangeStatement;
}

/**
* Same as {@link #getFindExpiredDurableSubRangeStatement()} for prioritized messages,
* where the ack table holds one last acked id per priority.
*/
public String getFindExpiredDurableSubRangeByPriorityStatement() {
if (findExpiredDurableSubRangeByPriorityStatement == null) {
findExpiredDurableSubRangeByPriorityStatement = "SELECT D.LAST_ACKED_ID, D.PRIORITY,"
+ " (SELECT MIN(B.ID) FROM " + getFullMessageTableName() + " B"
+ " WHERE B.CONTAINER=D.CONTAINER AND B.PRIORITY=D.PRIORITY AND B.ID > D.LAST_ACKED_ID"
+ " AND (B.XID IS NOT NULL OR B.EXPIRATION <= 0 OR B.EXPIRATION >= ?)),"
+ " (SELECT MAX(B.ID) FROM " + getFullMessageTableName() + " B"
+ " WHERE B.CONTAINER=D.CONTAINER AND B.PRIORITY=D.PRIORITY)"
+ " FROM " + getFullAckTableName() + " D"
+ " WHERE D.CONTAINER=? AND D.CLIENT_ID=? AND D.SUB_NAME=? AND D.XID IS NULL";
}
return findExpiredDurableSubRangeByPriorityStatement;
}

/**
* Expired messages with an id in a range, oldest first.
* Parameters: container, lower id (exclusive), upper id (exclusive), now.
*/
public String getFindExpiredMessagesInRangeStatement() {
if (findExpiredMessagesInRangeStatement == null) {
findExpiredMessagesInRangeStatement = "SELECT ID, MSG FROM " + getFullMessageTableName()
+ " WHERE CONTAINER=? AND ID > ? AND ID < ?"
+ " AND XID IS NULL AND EXPIRATION > 0 AND EXPIRATION < ?"
+ " ORDER BY ID";
}
return findExpiredMessagesInRangeStatement;
}

/**
* Same as {@link #getFindExpiredMessagesInRangeStatement()} limited to one priority.
* Parameters: container, lower id (exclusive), upper id (exclusive), now, priority.
*/
public String getFindExpiredMessagesInRangeByPriorityStatement() {
if (findExpiredMessagesInRangeByPriorityStatement == null) {
findExpiredMessagesInRangeByPriorityStatement = "SELECT ID, MSG FROM " + getFullMessageTableName()
+ " WHERE CONTAINER=? AND ID > ? AND ID < ?"
+ " AND XID IS NULL AND EXPIRATION > 0 AND EXPIRATION < ?"
+ " AND PRIORITY=?"
+ " ORDER BY ID";
}
return findExpiredMessagesInRangeByPriorityStatement;
}

public String getNextDurableSubscriberMessageStatement() {
if (nextDurableSubscriberMessageStatement == null) {
nextDurableSubscriberMessageStatement = "SELECT M.ID, M.MSG FROM "
Expand Down Expand Up @@ -884,6 +956,22 @@ public void setFindDurableSubMessagesStatement(String findDurableSubMessagesStat
this.findDurableSubMessagesStatement = findDurableSubMessagesStatement;
}

public void setFindExpiredDurableSubRangeStatement(String findExpiredDurableSubRangeStatement) {
this.findExpiredDurableSubRangeStatement = findExpiredDurableSubRangeStatement;
}

public void setFindExpiredDurableSubRangeByPriorityStatement(String findExpiredDurableSubRangeByPriorityStatement) {
this.findExpiredDurableSubRangeByPriorityStatement = findExpiredDurableSubRangeByPriorityStatement;
}

public void setFindExpiredMessagesInRangeStatement(String findExpiredMessagesInRangeStatement) {
this.findExpiredMessagesInRangeStatement = findExpiredMessagesInRangeStatement;
}

public void setFindExpiredMessagesInRangeByPriorityStatement(String findExpiredMessagesInRangeByPriorityStatement) {
this.findExpiredMessagesInRangeByPriorityStatement = findExpiredMessagesInRangeByPriorityStatement;
}

/**
* @param nextDurableSubscriberMessageStatement the nextDurableSubscriberMessageStatement to set
*/
Expand Down
Loading
Loading