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
