Skip to content

Implement recoverExpired() for the JDBC topic store (#2630) - #2634

Open
papinifrancesco wants to merge 4 commits into
apache:mainfrom
papinifrancesco:jdbc-topic-recover-expired
Open

papinifrancesco wants to merge 4 commits into
apache:mainfrom
papinifrancesco:jdbc-topic-recover-expired

Conversation

@papinifrancesco

Copy link
Copy Markdown

Link to #2630

Problems

  1. Topic expiry loads the whole JDBC backlog into the heap. Topic's expiry task uses TopicMessageStore.recoverExpired() only for KahaDB. Every other store falls back to doBrowse(), which on JDBC runs SELECT ID, MSG FROM ACTIVEMQ_MSGS WHERE CONTAINER=? ORDER BY ID with no row limit. Drivers that read the whole result set up front (PostgreSQL by default) load the entire durable backlog every expireMessagesPeriod, until the broker runs out of memory, even when no message has a TTL. Reproduction and CI evidence: https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom
  2. Expiring a message can drop earlier, non-expired messages of the same subscription. The JDBC store records a durable subscription's acks as a high-water mark (UPDATE ACTIVEMQ_ACKS SET LAST_ACKED_ID=?). When a message with a TTL expires behind one without a TTL, the expiry task (or a browse) acks the expired message, LAST_ACKED_ID moves past the earlier one, and after a restart that earlier message is never delivered. KahaDB tracks acks per message and is not affected.

Changes

  • JDBCTopicMessageStore.recoverExpired() is implemented. Per subscription, DefaultJDBCAdapter.doRecoverExpired() runs two plain-SQL statements (new in Statements, overridable like the others):

    1. per ack row: the last acked id, the id of the first pending message that is not expired or belongs to a prepared XA transaction, and the highest message id;
    2. the expired messages between the last acked id and that id, with setMaxRows(maxExpirePageSize).

    So only the run of expired messages directly after LAST_ACKED_ID is returned (per priority with prioritized messages), and acking them can never ack a message that has not expired. Messages already loaded for another subscription are shared, and the listener's hasSpace() is honoured as in the KahaDB implementation.

  • Topic uses recoverExpired() for JDBC as well as KahaDB. The getMessageCount() == 0 shortcut stays KahaDB-only, since on JDBC it is a COUNT query on every run. The recoverExpired() part of the expiry task moves to expireFromStore().

  • Topic.doBrowse() (JMX browse, StatisticsBroker) no longer acks the expired messages it browses on JDBC; it expires through expireFromStore() instead.

Behaviour change

On JDBC, expired messages queued behind a message that has not expired stay in the store until that message is consumed. They are still dropped at dispatch.

Tests

  • JDBCRecoverExpiredTest (H2): stopping at the first non-expired message, the page limit (messages shared between subscriptions), a subset of subscriptions, prioritized messages, a prepared XA transaction, the recovery listener.
  • JDBCDurableSubExpirationOrderTest (H2): end-to-end with a broker restart, so the result comes from the database and not from a cursor in memory, for both the expiry task and a topic browse. On current main without this change, testExpiredMessageDoesNotAckEarlierMessage and testBrowseDoesNotAckEarlierMessage fail.
  • JDBCPersistenceAdapterExpiredMessageTest.testMaxExpirePageSize now proxies recoverExpired(), which the expiry task calls instead of recover().
  • No regressions in the store/jdbc package and the topic expiry / durable subscription tests (DurableSubscriptionOffline*Test, ActiveDurableSubscriptionBrowseExpireTest, AMQ6122Test, JdbcDurableSubDupTest, MessageExpirationTest, ExpiredMessages*Test, KahaDBRecoverExpiredTest, ...). Four of those fail on my machine (Windows, 2 CPUs) with and without this change, so they are not caused by it: ExpiredMessagesTest (2), JDBCXACommitExceptionTest.testNonTxEnqueueOverNetworkErrorsRestart, JmsSendReceiveWithMessageExpirationTest.testConsumeExpiredTopicDurable.
  • On PostgreSQL 17 (https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom/actions/runs/37039551631): 6.3.2 and 6.2.10 with this change survive a 200,000 × 5 KB backlog of an offline durable subscriber that makes the unpatched brokers run out of memory. With a 30 s TTL on every message, the expiry task acks 400 messages per run with no out-of-memory error.

Left out

A manual browse of a JDBC topic still reads the whole result set on PostgreSQL (JDBCMessageStore.recover()). Only an explicit JMX or statistics browse triggers it, never the periodic task.

Disclosure

This change, its tests and the reproduction were written with an AI assistant (Claude Code, Anthropic). In line with the ASF guidance on generative tooling, the commits carry a Co-Authored-By trailer naming the tool.

🤖 Generated with Claude Code

papinifrancesco and others added 4 commits October 2, 2026 16:58
The topic expiry task fell back to Topic.doBrowse() for every store but
KahaDB. On JDBC that runs SELECT ID, MSG ... WHERE CONTAINER=? ORDER BY ID
with no row limit, and drivers that read the whole result set up front
(PostgreSQL by default) load the entire durable backlog into the heap every
expireMessagesPeriod, until the broker runs out of memory.

JDBCTopicMessageStore now implements recoverExpired() with one query per
subscription that returns only expired messages, at most
maxExpirePageSize of them.

The JDBC store acks a durable subscription by moving its LAST_ACKED_ID, so
expiring a message also acked every earlier message of the subscription,
including messages that had not expired: they were lost after a restart.
The query therefore only returns the expired messages that directly follow
LAST_ACKED_ID and stops at the first message that has not expired or is
part of a prepared XA transaction (per priority with prioritized messages).

Fixes apache#2630

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Topic.doBrowse() (JMX browse of a topic or durable subscription, the
StatisticsBroker plugin) expired every expired message it browsed. On JDBC
that ack moves the subscription's LAST_ACKED_ID, so it also acked the
earlier, non-expired messages, the same loss as the expiry task had. For
JDBC the browse now leaves the browsed messages alone and expires through
recoverExpired(), shared with the expiry task (expireFromStore()).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The single query evaluated its correlated MIN() subquery for every
candidate message, so one recoverExpired() call was quadratic in the
backlog when all of it was expired (Derby: 2.9 s for 1,000 messages,
93 s for 8,000; on PostgreSQL with 200,000 it never finished and the
expiry task, which skips runs while one is in progress, stopped).

doRecoverExpired() now reads, per ack row, the last acked id, the first
blocking message and the highest message id, then loads the expired
messages of that id range with the row limit. One call takes about
50-140 ms on Derby whatever the backlog.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Derby is being removed from ActiveMQ; the new tests use H2DB and
H2JDBCAdapter like the other H2 tests.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@jbonofre
jbonofre self-requested a review October 3, 2026 13:10
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant