From d579453fd16c688f5f2e69f3dda5490e3202c07a Mon Sep 17 00:00:00 2001 From: Francesco Papini <37743983+papinifrancesco@users.noreply.github.com> Date: Thu, 8 Oct 2026 20:17:33 +0200 Subject: [PATCH] Limit the topic browse query on the JDBC store (#2630) Topic.doBrowse(), also used by the topic expiry task for stores without recoverExpired(), recovered the topic through MessageStore.recover(). The JDBC store ran the query with no limit and relied on the listener to stop reading, but pgjdbc reads the whole result set in executeQuery(), so a large durable subscription backlog was loaded into the heap on every expiry run. Add MessageStore.recover(listener, maxReturned), whose default calls recover(listener), so stores that do not override it are unchanged. The JDBC store puts the limit on the query (setMaxRows). Topic.doBrowse() passes its page size. What a browse returns or expires does not change. Co-Authored-By: Claude Opus 5.5 --- .../apache/activemq/broker/region/Topic.java | 2 +- .../apache/activemq/store/MessageStore.java | 15 ++ .../activemq/store/ProxyMessageStore.java | 5 + .../store/ProxyTopicMessageStore.java | 5 + .../activemq/store/jdbc/JDBCAdapter.java | 2 + .../activemq/store/jdbc/JDBCMessageStore.java | 11 +- .../jdbc/adapter/DefaultJDBCAdapter.java | 14 ++ ...CPersistenceAdapterExpiredMessageTest.java | 4 +- .../store/jdbc/JDBCTopicBrowseLimitTest.java | 179 ++++++++++++++++++ 9 files changed, 233 insertions(+), 4 deletions(-) create mode 100644 activemq-unit-tests/src/test/java/org/apache/activemq/store/jdbc/JDBCTopicBrowseLimitTest.java 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; + } + } +}