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


Reply via email to