Andrea Cosentino created CAMEL-25171:
----------------------------------------
Summary: camel-couchbase: the consumer discards route failures,
and with consumerProcessedStrategy=delete the document is already gone
Key: CAMEL-25171
URL: https://issues.apache.org/jira/browse/CAMEL-25171
Project: Camel
Issue Type: Bug
Components: camel-couchbase
Reporter: Andrea Cosentino
Assignee: Andrea Cosentino
{{CouchbaseConsumer.processBatch}} hands each exchange to the route and then
ignores what happened to it:
{code:java}
for (int index = 0; index < total && this.isBatchAllowed(); ++index) {
Exchange exchange = (Exchange) exchanges.poll();
exchange.setProperty(ExchangePropertyKey.BATCH_INDEX, index);
exchange.setProperty(ExchangePropertyKey.BATCH_SIZE, total);
exchange.setProperty(ExchangePropertyKey.BATCH_COMPLETE, index == total -
1);
this.pendingExchanges = total - index - 1;
getProcessor().process(exchange);
}
{code}
A failing route does not throw out of {{process()}} - the consumer processor is
async, so the failure is left on {{exchange.getException()}}. Nothing here
reads it. {{getExceptionHandler()}} does not appear anywhere in
{{CouchbaseConsumer}} at all, so {{bridgeErrorHandler}} is inert for this
component and a failed exchange produces no log line whatsoever.
This is the same defect as CAMEL-25024 in camel-mongodb.
h2. Why it is worse here: the document has already been deleted
With {{consumerProcessedStrategy=delete}}, the removal happens in the polling
loop - in {{pollWithSqlQuery}} and {{pollWithView}}, while the exchanges are
being built - and therefore *before* {{processBatch}} runs:
{code:java}
if ("delete".equalsIgnoreCase(consumerProcessedStrategy)) {
CouchbaseCollectionOperation.removeDocument(collection, id,
endpoint.getWriteQueryTimeout(),
endpoint.getConsumerRetryPause());
}
{code}
So when the route fails, the document is already gone from Couchbase, the
failure is swallowed, and nothing is logged. The message is lost silently.
h2. The batch bookkeeping was copied without its safety net
{{processBatch}} is a copy of {{GenericFileConsumer.processBatch}} with the
recovery machinery removed. The original keeps a {{notStarted}} queue,
decrements {{answer}} for exchanges that did not start, and drains *both*
{{notStarted}} and the leftover {{exchanges}} queue before returning. The
couchbase copy keeps the {{int answer = total;}} shape but never adjusts it,
and never drains the remainder - so exchanges still queued when
{{isBatchAllowed()}} flips to false at shutdown are neither processed nor
released. With a pooled exchange factory those are taken from the pool and
never returned.
h2. Proposed fix
Introduce the {{processExchange}} shape used in camel-mongodb: check
{{exchange.getException()}} after {{process()}}, report it through
{{getExceptionHandler()}}, and stop treating a failed exchange as delivered.
Release the exchanges left in the queue when the batch is cut short.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)