[
https://issues.apache.org/jira/browse/CAMEL-24933?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Work on CAMEL-24933 started by Andrea Cosentino.
------------------------------------------------
> camel-pulsar - the consumer leaks a pooled exchange for every message
> ---------------------------------------------------------------------
>
> Key: CAMEL-24933
> URL: https://issues.apache.org/jira/browse/CAMEL-24933
> Project: Camel
> Issue Type: Bug
> Components: camel-pulsar
> Reporter: Andrea Cosentino
> Assignee: Andrea Cosentino
> Priority: Major
>
> h3. Summary
> The consumer creates an {{Exchange}} through the pooled exchange factory and
> then throws it away for a
> copy. With {{camel.main.exchange-factory=pooled}} every consumed message
> leaks one pooled exchange and
> makes {{PooledExchangeFactory.release}} throw a {{ClassCastException}} that
> is swallowed at DEBUG.
> h3. Details
> {{PulsarMessageListener.received()}}:
> {code:java}
> final Exchange exchange = PulsarMessageUtils.updateExchange(message,
> pulsarConsumer.createExchange(false));
> {code}
> and {{PulsarMessageUtils.updateExchange}} starts with:
> {code:java}
> final Exchange output = input.copy();
> {code}
> {{DefaultPooledExchange.newCopy()}} returns {{new DefaultExchange(this)}},
> i.e. a copy is never pooled.
> So with pooling enabled:
> # {{createExchange(false)}} takes a {{DefaultPooledExchange}} out of the
> factory;
> # {{copy()}} discards it in favour of a plain {{DefaultExchange}} - the
> pooled one is never {{done()}}
> and never returned to the pool;
> # the callback calls {{releaseExchange(copy, false)}}, and since the copy is
> not a {{PooledExchange}}
> the {{pooledExchange.done()}} branch in {{DefaultConsumer.releaseExchange}}
> is skipped;
> # {{PooledExchangeFactory.release}} opens with {{PooledExchange ee =
> (PooledExchange) exchange;}}, so it
> throws {{ClassCastException}}, catches it, logs {{"Error resetting
> exchange: ..."}} at DEBUG and
> increments the {{discarded}} statistic.
> The result is one orphaned pooled exchange and one thrown-and-swallowed
> {{ClassCastException}} per
> consumed message. The pool never refills, so pooling brings this component no
> benefit and some cost.
> Without pooling the copy is "only" a second exchange allocation per message.
> h3. Proposed fix
> Populate the exchange the consumer created instead of copying it.
> {{updateExchange}} has exactly one
> caller, so the change is contained. Note also that {{output.setIn(msg)}} in
> that method sets back the
> message it just read, and that the sibling {{updateExchangeWithException}} is
> dead code - nothing calls
> it.
> ----
> _Reported by Claude Code on behalf of oscerd (Andrea Cosentino)._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)