papinifrancesco opened a new issue, #2630: URL: https://github.com/apache/activemq/issues/2630
### Summary With the JDBC persistence adapter, the periodic topic expiry task (`Topic.expireMessagesTask`, every `expireMessagesPeriod`, default 30 s) falls back to `Topic.doBrowse()`. That calls `JDBCMessageStore.recover()` → `DefaultJDBCAdapter.doRecover()`, which runs ```sql SELECT ID, MSG FROM ACTIVEMQ_MSGS WHERE CONTAINER=? ORDER BY ID ``` with no `setMaxRows()` and no fetch size. [AMQ-6067](https://issues.apache.org/jira/browse/AMQ-6067) made the recovery listener stop after `maxExpirePageSize` (400) messages. That is enough for drivers that stream rows, but the PostgreSQL driver reads the **entire** result set into memory inside `executeQuery()`, before the first `rs.next()`. On PostgreSQL, every run therefore loads every stored message of the topic (the whole durable-subscription backlog) into the heap, only to use 400 of them. **Reproduction:** https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom reproduces it on 6.2.10 and 6.3.2, with synthetic data, heap dumps and a one-click GitHub Actions workflow (details below). ### Environment - ActiveMQ 6.2.6 when it happened; the same code is in 6.2.10 and 6.3.2 (references below) - JDBC persistence adapter on PostgreSQL (pgjdbc 42.7.x, commons-dbcp2), default statements and table names - One topic with a durable subscription whose consumer fell behind during a bulk load: ~347,000 persistent messages of ~5 KB stored - JVM heap ≈ 2.4 GB, default `expireMessagesPeriod`, no message sets a TTL ### What happened The broker stopped with `java.lang.OutOfMemoryError: Java heap space`. In the heap dump (Eclipse MAT, Leak Suspects), **69% of the heap (1.78 GB)** was one `ArrayList` local to the thread `ActiveMQ BrokerService[activemq] Task-614`, holding **346,764 `org.postgresql.core.Tuple`** (693,531 `byte[]`). Stack of that thread (line numbers from 6.2.6): ``` org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2387) org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:372) org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:517) org.postgresql.jdbc.PgPreparedStatement.executeQuery(PgPreparedStatement.java:137) org.apache.commons.dbcp2.DelegatingPreparedStatement.executeQuery(DelegatingPreparedStatement.java:123) org.apache.activemq.store.jdbc.adapter.DefaultJDBCAdapter.doRecover(DefaultJDBCAdapter.java:406) org.apache.activemq.store.jdbc.JDBCMessageStore.recover(JDBCMessageStore.java:279) org.apache.activemq.store.ProxyTopicMessageStore.recover(ProxyTopicMessageStore.java:63) org.apache.activemq.broker.region.Topic.doBrowse(Topic.java:678) org.apache.activemq.broker.region.Topic.lambda$new$2(Topic.java:920) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ``` The thread is still inside `executeQuery()`, receiving tuples: the listener's `hasSpace()` check (AMQ-6067) never gets a chance to stop it. ### Reproduction (synthetic, public) https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom uses only the official binary distributions, a stock PostgreSQL 17 and the broker's own `activemq consumer` / `activemq producer` CLI. A GitHub Actions workflow runs four cases. Each run offers 200,000 × 5 KB persistent messages, no TTL, to a topic with one offline durable subscriber, on a 1 GB heap: | ActiveMQ | Store | `expireMessagesPeriod` | Result | |---|---|---|---| | 6.2.10 | JDBC / PostgreSQL | default | `OutOfMemoryError` at ~86k stored rows, `Failed to browse Topic` WARN, browse caught 3× in `DefaultJDBCAdapter.doRecover` | | 6.3.2 | JDBC / PostgreSQL | default | `OutOfMemoryError` at ~89k stored rows, broker exits, same stack | | 6.2.10 | JDBC / PostgreSQL | `0` | all 200,000 rows stored, no OOM | | 6.2.10 | KahaDB | default | no OOM | - Run: https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom/actions/runs/36845523696 - Evidence, including both heap dumps and a MAT Leak Suspects report: https://github.com/papinifrancesco/activemq-jdbc-topic-expiry-oom/releases/tag/evidence-2026-10-01 MAT on the 6.2.10 dump: thread `ActiveMQ BrokerService[repro] Task-3` retains 53.8% of the heap in 54,343 `org.postgresql.core.Tuple` objects. Its stack is the same as above (`Topic.lambda$new$2:950` → `Topic.doBrowse:691` → `JDBCMessageStore.recover:279` → `DefaultJDBCAdapter.doRecover:406` → `executeQuery`). Most of the remaining heap (28%) was the durable subscription's pending cache. We capped that separately with a topic `memoryLimit`; it is not part of this report. ### Code references (6.3.2; same code on `activemq-6.2.x`) - `Topic.java:930` uses `store.recoverExpired(...)` only when `store.getType() == StoreType.KAHADB` (added by [AMQ-9698](https://issues.apache.org/jira/browse/AMQ-9698)). Any other store takes `Topic.java:980-982`: `doBrowse(new InsertionCountList<>(), getMaxExpirePageSize())`. On `activemq-6.2.x` these are lines 898 and 950. - `DefaultJDBCAdapter.java:399-404`: `doRecover` runs `findAllMessagesStatement` with no limit. The other recovery queries in the same class do limit it: `setMaxRows(...)` at lines 433, 593, 627 and 1068. - `JDBCMessageStore.java:274-296`: the `hasSpace()` check per row (AMQ-6067) runs after `executeQuery()` has returned. - Both JDBC files have the same line numbers on `activemq-6.2.x`. ### Why the AMQ-6067 fix does not cover PostgreSQL AMQ-6067 was reported on Oracle, whose driver streams rows (default fetch size 10), so stopping after 400 rows fixed it there. pgjdbc buffers the whole result set unless the connection is **not** in autocommit mode **and** a fetch size is set ([pgjdbc docs, "Getting results based on a cursor"](https://jdbc.postgresql.org/documentation/query/)). `doRecover` sets no fetch size, and the stack above shows the driver buffering the full result set. MySQL Connector/J also buffers full result sets by default, so it is probably affected too; we have not tested it. ### Expected behaviour The expiry browse of a JDBC topic store reads at most `maxExpirePageSize` rows, as it effectively does with KahaDB. ### Possible fixes (for discussion) 1. Give `doRecover` a limit and call `s.setMaxRows(limit)`, like the other recovery methods. `doBrowse` already knows `max`. 2. Implement `recoverExpired` for the JDBC store with a bounded query on the existing indexed column (`ACTIVEMQ_MSGS_EIDX`), e.g. `... WHERE CONTAINER=? AND EXPIRATION > 0 AND EXPIRATION < ?`. The fallback would then be unnecessary, and a topic whose messages have no TTL would return zero rows. 3. At least, skip the browse when no stored message of the topic has `EXPIRATION > 0`. ### Workaround `expireMessagesPeriod="0"` on the topic `policyEntry`. Expired messages are still discarded at dispatch (`PrefetchSubscription.dispatchPending`). They are no longer removed while a durable subscriber is offline, though: the JDBC cleanup (`doDeleteOldMessages`) only deletes acknowledged rows. The production heap dump cannot be shared, because it contains customer data. The reproduction above gives the same stack and the same retention pattern with synthetic data only. -- 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
