papinifrancesco opened a new pull request, #2634:
URL: https://github.com/apache/activemq/pull/2634

   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](https://claude.com/claude-code)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
For further information, visit: https://activemq.apache.org/contact


Reply via email to