diff --git a/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java b/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java index c6860d1c5ce..f6806403e1f 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java +++ b/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java @@ -744,7 +744,7 @@ public boolean hasSpace() { public boolean isDuplicate(MessageId id) { return false; } - }); + }, max); final ConnectionContext connectionContext = createConnectionContext(); for (Message message : toExpire) { for (DurableTopicSubscription sub : durableSubscribers.values()) { diff --git a/activemq-broker/src/main/java/org/apache/activemq/store/MessageStore.java b/activemq-broker/src/main/java/org/apache/activemq/store/MessageStore.java index 70945fc6173..e297d57b52c 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/store/MessageStore.java +++ b/activemq-broker/src/main/java/org/apache/activemq/store/MessageStore.java @@ -138,6 +138,21 @@ public interface MessageStore extends Service { */ void recover(MessageRecoveryListener container) throws Exception; + /** + * Recover at most maxReturned messages to be delivered. + * + * The default implementation relies on {@link MessageRecoveryListener#hasSpace()} to + * stop the recovery. A store that would otherwise read every message before the + * listener can stop it should bound the read itself. + * + * @param listener + * @param maxReturned the maximum number of messages to recover + * @throws Exception + */ + default void recover(MessageRecoveryListener listener, int maxReturned) throws Exception { + recover(listener); + } + /** * The destination that the message store is holding messages for. * diff --git a/activemq-broker/src/main/java/org/apache/activemq/store/ProxyMessageStore.java b/activemq-broker/src/main/java/org/apache/activemq/store/ProxyMessageStore.java index dfb369bd18f..7274d040e72 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/store/ProxyMessageStore.java +++ b/activemq-broker/src/main/java/org/apache/activemq/store/ProxyMessageStore.java @@ -60,6 +60,11 @@ public void recover(MessageRecoveryListener listener) throws Exception { delegate.recover(listener); } + @Override + public void recover(MessageRecoveryListener listener, int maxReturned) throws Exception { + delegate.recover(listener, maxReturned); + } + @Override public void removeAllMessages(ConnectionContext context) throws IOException { delegate.removeAllMessages(context); diff --git a/activemq-broker/src/main/java/org/apache/activemq/store/ProxyTopicMessageStore.java b/activemq-broker/src/main/java/org/apache/activemq/store/ProxyTopicMessageStore.java index d9b92500c01..07d3c19098a 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/store/ProxyTopicMessageStore.java +++ b/activemq-broker/src/main/java/org/apache/activemq/store/ProxyTopicMessageStore.java @@ -63,6 +63,11 @@ public void recover(MessageRecoveryListener listener) throws Exception { delegate.recover(listener); } + @Override + public void recover(MessageRecoveryListener listener, int maxReturned) throws Exception { + delegate.recover(listener, maxReturned); + } + @Override public void removeAllMessages(ConnectionContext context) throws IOException { delegate.removeAllMessages(context); diff --git a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCAdapter.java b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCAdapter.java index dc57b24e3b0..95cba915e13 100644 --- a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCAdapter.java +++ b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCAdapter.java @@ -52,6 +52,8 @@ public interface JDBCAdapter { void doRecover(TransactionContext c, ActiveMQDestination destination, JDBCMessageRecoveryListener listener) throws Exception; + void doRecover(TransactionContext c, ActiveMQDestination destination, int maxReturned, JDBCMessageRecoveryListener listener) throws Exception; + void doSetLastAck(TransactionContext c, ActiveMQDestination destination, XATransactionId xid, String clientId, String subscriptionName, long seq, long prio) throws SQLException, IOException; void doRecoverSubscription(TransactionContext c, ActiveMQDestination destination, String clientId, String subscriptionName, JDBCMessageRecoveryListener listener) diff --git a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCMessageStore.java b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCMessageStore.java index 78c27b71839..c7dcdb1cd42 100644 --- a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCMessageStore.java +++ b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/JDBCMessageStore.java @@ -272,11 +272,20 @@ public void removeMessage(ConnectionContext context, MessageAck ack) throws IOEx @Override public void recover(final MessageRecoveryListener listener) throws Exception { + recover(listener, 0); + } + + /** + * The limit is applied to the query: the listener alone cannot bound the memory used + * with drivers that read the whole result set up front. + */ + @Override + public void recover(final MessageRecoveryListener listener, int maxReturned) throws Exception { // Get all the Message ids out of the database. TransactionContext c = persistenceAdapter.getTransactionContext(); try { - adapter.doRecover(c, destination, new JDBCMessageRecoveryListener() { + adapter.doRecover(c, destination, maxReturned, new JDBCMessageRecoveryListener() { @Override public boolean recoverMessage(long sequenceId, byte[] data) throws Exception { if (listener.hasSpace()) { diff --git a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/adapter/DefaultJDBCAdapter.java b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/adapter/DefaultJDBCAdapter.java index 47823ee649a..7ddc02c3684 100644 --- a/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/adapter/DefaultJDBCAdapter.java +++ b/activemq-jdbc-store/src/main/java/org/apache/activemq/store/jdbc/adapter/DefaultJDBCAdapter.java @@ -398,10 +398,24 @@ public void doRemoveMessage(TransactionContext c, long seq, XATransactionId xid) @Override public void doRecover(TransactionContext c, ActiveMQDestination destination, JDBCMessageRecoveryListener listener) throws Exception { + doRecover(c, destination, 0, listener); + } + + /** + * @param maxReturned the maximum number of rows to read, or 0 for no limit. Some drivers + * (PostgreSQL by default) read the whole result set in executeQuery(), + * so stopping the listener early does not bound the memory used. + */ + @Override + public void doRecover(TransactionContext c, ActiveMQDestination destination, int maxReturned, + JDBCMessageRecoveryListener listener) throws Exception { PreparedStatement s = null; ResultSet rs = null; try { s = c.getConnection().prepareStatement(this.statements.getFindAllMessagesStatement()); + if (maxReturned > 0) { + s.setMaxRows(maxReturned); + } s.setString(1, destination.getQualifiedName()); rs = s.executeQuery(); if (this.statements.isUseExternalMessageReferences()) { diff --git a/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCPersistenceAdapterExpiredMessageTest.java b/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCPersistenceAdapterExpiredMessageTest.java index e470dd70ffe..33c782b2ffa 100644 --- a/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCPersistenceAdapterExpiredMessageTest.java +++ b/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCPersistenceAdapterExpiredMessageTest.java @@ -75,7 +75,7 @@ public TopicMessageStore createTopicMessageStore(ActiveMQTopic destination) thro ProxyTopicMessageStore proxy = new ProxyTopicMessageStore(super.createTopicMessageStore(destination)) { @Override - public void recover(final MessageRecoveryListener listener) throws Exception { + public void recover(final MessageRecoveryListener listener, int maxReturned) throws Exception { MessageRecoveryListener delegate = new MessageRecoveryListener() { @Override @@ -99,7 +99,7 @@ public boolean hasSpace() { return listener.hasSpace(); } }; - super.recover(delegate); + super.recover(delegate, maxReturned); } }; diff --git a/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCTopicBrowseLimitTest.java b/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCTopicBrowseLimitTest.java new file mode 100644 index 00000000000..ca37277a309 --- /dev/null +++ b/activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCTopicBrowseLimitTest.java @@ -0,0 +1,179 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.activemq.store.jdbc; + +import static org.junit.Assert.assertEquals; + +import java.io.File; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import jakarta.jms.Connection; +import jakarta.jms.MessageProducer; +import jakarta.jms.Session; + +import org.apache.activemq.ActiveMQConnectionFactory; +import org.apache.activemq.broker.BrokerService; +import org.apache.activemq.broker.region.Destination; +import org.apache.activemq.broker.region.policy.PolicyEntry; +import org.apache.activemq.broker.region.policy.PolicyMap; +import org.apache.activemq.command.ActiveMQDestination; +import org.apache.activemq.command.ActiveMQTopic; +import org.apache.activemq.command.Message; +import org.apache.activemq.command.MessageId; +import org.apache.activemq.store.MessageRecoveryListener; +import org.apache.activemq.store.MessageStore; +import org.apache.activemq.store.jdbc.adapter.H2JDBCAdapter; +import org.apache.activemq.store.jdbc.h2.H2DB; +import org.apache.activemq.test.annotations.ParallelTest; +import org.junit.After; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.experimental.categories.Category; +import org.junit.rules.TemporaryFolder; +import org.junit.rules.Timeout; + +/** + * A topic browse, also used by the topic expiry task for the JDBC store, must not read more + * rows than the browse page size: some drivers (PostgreSQL by default) read the whole result + * set in executeQuery(), so a large durable subscription backlog could exhaust the heap. + */ +@Category(ParallelTest.class) +public class JDBCTopicBrowseLimitTest { + + private static final int MESSAGE_COUNT = 50; + private static final int PAGE_SIZE = 10; + + @Rule + public Timeout globalTimeout = new Timeout(60, TimeUnit.SECONDS); + + @Rule + public TemporaryFolder dataFileDir = new TemporaryFolder(new File("target")); + + private final ActiveMQTopic topic = new ActiveMQTopic("test.topic"); + private final List recoverLimits = new CopyOnWriteArrayList<>(); + private BrokerService broker; + private Connection connection; + + @Before + public void startBroker() throws Exception { + broker = new BrokerService(); + broker.setUseJmx(false); + broker.setSchedulerSupport(false); + broker.setDataDirectoryFile(dataFileDir.getRoot()); + JDBCPersistenceAdapter jdbc = new JDBCPersistenceAdapter(); + jdbc.setDataSource(H2DB.createDataSource("JDBCTopicBrowseLimitTest")); + // records the row limit of each recover query + jdbc.setAdapter(new H2JDBCAdapter() { + @Override + public void doRecover(TransactionContext c, ActiveMQDestination destination, int maxReturned, + JDBCMessageRecoveryListener listener) throws Exception { + recoverLimits.add(maxReturned); + super.doRecover(c, destination, maxReturned, listener); + } + }); + broker.setPersistenceAdapter(jdbc); + broker.setDeleteAllMessagesOnStartup(true); + PolicyMap policyMap = new PolicyMap(); + PolicyEntry policy = new PolicyEntry(); + // the test drives the browse itself + policy.setExpireMessagesPeriod(0); + policy.setMaxBrowsePageSize(PAGE_SIZE); + policyMap.setDefaultEntry(policy); + broker.setDestinationPolicy(policyMap); + broker.start(); + broker.waitUntilStarted(); + + connection = new ActiveMQConnectionFactory("vm://localhost").createConnection(); + connection.setClientID("clientId"); + connection.start(); + Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + // an offline durable subscription keeps the messages in the store + session.createDurableSubscriber(topic, "sub1").close(); + MessageProducer producer = session.createProducer(topic); + for (int i = 0; i < MESSAGE_COUNT; i++) { + producer.send(session.createTextMessage("message" + i)); + } + session.close(); + } + + @After + public void stopBroker() throws Exception { + if (connection != null) { + connection.close(); + } + if (broker != null) { + broker.stop(); + broker.waitUntilStopped(); + } + } + + @Test + public void testRecoverReadsAtMostMaxReturnedRows() throws Exception { + MessageStore store = broker.getDestination(topic).getMessageStore(); + + // a listener that never stops the recovery: only the query can bound it + AtomicInteger count = new AtomicInteger(); + store.recover(new CountingListener(count), PAGE_SIZE); + assertEquals(PAGE_SIZE, count.get()); + + count.set(0); + store.recover(new CountingListener(count)); + assertEquals(MESSAGE_COUNT, count.get()); + } + + @Test + public void testBrowseLimitsTheQuery() throws Exception { + Destination destination = broker.getDestination(topic); + recoverLimits.clear(); + + assertEquals(PAGE_SIZE, destination.browse().length); + assertEquals(List.of(PAGE_SIZE), recoverLimits); + } + + private static class CountingListener implements MessageRecoveryListener { + private final AtomicInteger count; + + CountingListener(AtomicInteger count) { + this.count = count; + } + + @Override + public boolean recoverMessage(Message message) { + count.incrementAndGet(); + return true; + } + + @Override + public boolean recoverMessageReference(MessageId ref) { + return true; + } + + @Override + public boolean hasSpace() { + return true; + } + + @Override + public boolean isDuplicate(MessageId ref) { + return false; + } + } +}