Repository navigation
CAMEL-25221: camel-couchbase - remove a document only after its exchange completed - #27347
Conversation
…nge completed With consumerProcessedStrategy=delete the consumer removed every document of a poll while building the exchanges, before processBatch handed any of them to the route. A document was therefore lost when its exchange failed, and also when it was never delivered because the batch was cut short by maxMessagesPerPoll or by the consumer stopping. The removal now runs in an on-completion of the exchange, as the delete-after-read consumers of other components do, so it happens only for an exchange that went through the route and completed. A failed or undelivered document is kept for a later poll. A failure to remove the document is reported through the consumer's exception handler. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
🌟 Thank you for your contribution to the Apache Camel project! 🌟 🐫 Apache Camel Committers, please review the following items:
|
davsclaus
left a comment
There was a problem hiding this comment.
Thanks, this fixes the data loss raised on #27123: documents of failed exchanges, of rows beyond maxMessagesPerPoll, and of rows dropped at shutdown were deleted. Deleting on completion follows the established pattern (as in aws2-s3), and the upgrade-guide text and tests cover the new behaviour. One question inline.
Claude Code on behalf of davsclaus. This review was generated by an AI agent and may contain inaccuracies. Please verify all suggestions before applying. It does not replace specialized review tools or static analysis.
|
🧪 CI tested the following changed modules:
🔬 Scalpel shadow comparison — Scalpel: 9 of 698 tested, 26 compile-only — current: 9 all testedMaveniverse Scalpel detected 9 affected modules (current approach: 9). Skip-tests mode would test 9 modules (2 direct + 8 downstream), skip tests for 26 (generated code, meta-modules) Modules Scalpel would test (9)
Modules with tests skipped (26)
All tested modules (36 modules, 5m 28s total)Total reactor time: 5m 28s
Top 20 slowest modules:
|
…its exchange is in flight With consumerProcessedStrategy=delete the document is removed in an on-completion that is handed over with the exchange, so it is not removed before an asynchronous part of the route has processed it. When the route hands the exchange over to another thread (seda with waitForTaskToComplete=Never, an asynchronous producer), the next poll could read the same document again before that exchange completed and deliver it twice. The consumer now keeps the IDs of the documents whose exchange is in flight, as the in-progress repository of the file and aws2-s3 consumers does, and skips them in later polls. An ID is released when its exchange completes or fails, when its exchange is never handed to the route (batch cut short by maxMessagesPerPoll or by the consumer stopping), and when the poll fails part way. A skipped row does not count towards maxMessagesPerPoll. The other strategies leave the document in place and are not guarded. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
gnodet-bot
left a comment
There was a problem hiding this comment.
Solid fix for the data loss in CAMEL-25221. The delete-on-completion pattern correctly follows the established Camel conventions (aws2-s3, file consumer), and the in-flight guard handles the async handover case cleanly.
Verified:
- Thread safety:
ConcurrentHashMap.newKeySet()is the right choice — poll-side writes are under the existinglock, callback-sideremove()is thread-safe and idempotent. The check-then-act inisInFlight/inFlight.adddoesn't race because both run under the poll lock. - Resource cleanup:
releaseUndeliveredcorrectly drains in-flight IDs before releasing pooled exchanges, so even ifPooledExchange.done()fires on-completion callbacks, the double-remove is harmless. - Failure partway through a poll: If
getDocument()throws for doc N, docs 0..N-1 are already in theexchangesqueue with their IDs ininFlight. The catch inpoll()→releaseUndeliveredcleans up both correctly. processBatchreturn value: Returns the number of exchanges actually built (excluding skipped in-flight docs), which is the correct semantics for the framework.- Test coverage: 15 tests exercise both query paths (SQL++, View), success/failure, maxMessagesPerPoll, shutdown mid-batch, async handover with completion/failure, partial poll failure, and control cases for other strategies. Thorough.
- Upgrade guide: Accurate, covers the behavior change including the edge case of per-consumer scope and the interaction with error handling (
onException, dead letter channel).
ast-grep flagged 3 broad-exception-catch patterns — all are standard Camel patterns (catch-rethrow in poll(), catch-handle in onComplete, pre-existing processExchange). Not issues.
This review was generated by an AI agent, Hermès on behalf of @gnodet.
davsclaus
left a comment
There was a problem hiding this comment.
Re-approving after the in-flight tracking commit; the set is only checked/added under the poll lock and is cleared on completion, failure and early stop.
Follow-up thought: if a handed-over exchange never completes (seda queue full/purged, async producer that never calls back), its id stays in the set for the life of the consumer and the document is skipped forever, even across stopRoute/startRoute since the consumer is reused. Clearing the set in doStop/doStart would bound that. Also worth a line in the upgrade note: a view emitting several rows for one document now delivers only the first row under delete.
|
Thanks. Good points, I'll take both in a follow-up once this is merged: clear the in-flight set in Claude Code on behalf of allthingssecurity |
oscerd
left a comment
There was a problem hiding this comment.
Traced the whole in-flight lifecycle; correct, and it's the real fix CAMEL-25221 asked for (the delete-during-poll was losing documents on failure and on the rows past maxMessagesPerPoll/shutdown).
- Delete after successful processing. The removal is an
onCompleteon-completion added when the exchange is built, so it runs only for an exchange that went through the route's unit of work and completed;onFailurekeeps the document for the next poll, and a handled exception (handled(true)/DLC) completes the exchange so the document is removed — matching the file/aws2-s3 delete-after-read pattern. Keeping the defaultallowHandover()==truemeans a handed-over exchange (seda waitForTaskToComplete=Never, async producer) deletes the document when that exchange completes, not when the consumer's unit of work ends, which is the data loss being fixed. - The in-flight guard is released on every path, so no document is leaked into a permanent skip:
onComplete(remove +finallyrelease),onFailure(release, keep doc),processBatch's trailingreleaseUndeliveredfor rows cut off bymaxMessagesPerPollor a stopping consumer, and the poll's innercatch → releaseUndeliveredwhen building the batch throws part way. A still-in-flight id is skipped on later polls (isInFlight) without counting towardmaxMessagesPerPoll, and the set is per-consumer so it doesn't mask another node.none/filterare untouched. - A removal failure now goes through
getExceptionHandler()instead of throwing out ofpoll()and aborting the rest of the batch.
The upgrade-guide subsection correctly documents the new timing, the "a route that always fails re-reads the document each poll" consequence, and the per-consumer in-flight limit. CouchbaseConsumerProcessedStrategyTest covers doc-present-during-route/removed-after, failed-keeps-doc (SQL++ and view), maxMessagesPerPoll=1, deferShutdown(CompleteCurrentTaskOnly), the removal-failure-to-exception-handler path, the none control, and the handover guard via seda:async. CI is green.
Approving.
This review was generated with AI assistance and reviewed/issued by the human operator. Claude Code on behalf of oscerd
…ument the fullDocument default With useView=true and fullDocument=false the consumer set the message body to row.valueAs(Object.class), which the Couchbase SDK 3 returns as an Optional, so the body was Optional[value] (Optional.empty when the view emitted null) instead of the value. The body is now the value, or null when the view emitted none. The SQL++ path is not affected: it uses the query row as a JSON string. The fullDocument option was documented with defaultValue false since it was added (CAMEL-15792), but the field has always defaulted to true, which keeps the behaviour from before the option existed (the consumer always fetched the document). The annotation now says true, and the component JSON, the catalog and the endpoint DSL are regenerated. No change at runtime. Upgrade guide: the body change, and, as asked in the review of apache#27347 (CAMEL-25221), that a view or SQL++ query returning several rows for one document delivers only the first of them with consumerProcessedStrategy=delete. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Description
CAMEL-25221
With
consumerProcessedStrategy=delete,CouchbaseConsumerremoved every document of a poll inpollWithSqlQuery/pollWithView, while building the exchanges and beforeprocessBatchhanded any of them to the route. So a document was lost:processBatchstops atmaxMessagesPerPolland whenisBatchAllowed()turns false because the consumer is stopping, and the rows after that point had already been removed.@davsclaus raised this in the review of #27123 (CAMEL-25171), and @oscerd filed it as CAMEL-25221 with the follow-up options. This PR takes the first one, delete after the exchange has been processed successfully, which the ticket calls the real fix (the comparable camel-mongodb change was CAMEL-25024).
This change:
SynchronizationAdapter.onComplete) added to the exchange when it is built, the delete-after-read pattern of other consumers (for example aws2-s3). It runs only for an exchange that went through the route's unit of work and completed: a failed exchange keeps its document for the next poll, and an exchange that was never handed to the route (released by the drain loop from CAMEL-25169, CAMEL-25170, CAMEL-25171, CAMEL-25172: camel-couchbase - close the connection, apply connectTimeout, report route failures, accept persistTo=2 #27123) never runs it. A handled exception (onException().handled(true), dead letter channel) completes the exchange, so the document is removed, as for the file consumer;allowHandover() == true, so when the route hands the exchange over to another thread (sedawithwaitForTaskToComplete=Never, an asynchronous producer) the document is removed when that exchange completes, not before the asynchronous part has processed it (withfalseit would be removed when the consumer's own unit of work ends, which is the data loss this PR fixes). Because that can be after the next poll, the consumer now keeps the IDs of the documents whose exchange is in flight and skips them in later polls, as the in-progress repository of the file and aws2-s3 consumers does (second commit, from @davsclaus's review). An ID is released when its exchange completes or fails, when the exchange is never handed to the route (batch cut short bymaxMessagesPerPollor by the consumer stopping), and when the poll fails part way (for example thefullDocumentget of a later row throws). A skipped row creates no exchange, so it does not count towardsmaxMessagesPerPoll. It is an internal concurrent set rather than a pluggableinProgressRepositoryendpoint option like aws2-s3's: the guard only has to cover this consumer's own handed-over exchanges, an unbounded set cannot evict an ID that is still in flight (aws2-s3's default is a 10000-entry FIFO cache), and it adds no new option to a bug fix. It is per consumer, so it does not stop a consumer on another node from reading the same document. It applies todeleteonly: withnone/filterthe document stays in the bucket and later polls read it again, as before;poll()and abort the rest of the batch);processBatch/processExchange(added by CAMEL-25169, CAMEL-25170, CAMEL-25171, CAMEL-25172: camel-couchbase - close the connection, apply connectTimeout, report route failures, accept persistTo=2 #27123) are updated;filterandnoneare unchanged (neither removes anything);=== camel-couchbaseentry (from CAMEL-25169, CAMEL-25170, CAMEL-25171, CAMEL-25172: camel-couchbase - close the connection, apply connectTimeout, report route failures, accept persistTo=2 #27123) said the document "is removed during the poll, before the route runs"; it is replaced by a==== consumerProcessedStrategy=deletesubsection that describes the new timing, including that a document whose route always fails is now consumed again on every poll, and the in-flight guard with its per-consumer limit.Tests: new
CouchbaseConsumerProcessedStrategyTestruns the consumer in a real route (so the exchange gets a unit of work) with the scheduler not started, callspoll()against a mocked cluster, and verifiesCollection.remove:maxMessagesPerPoll=1with three rows: only the first is removed;deferShutdown(CompleteCurrentTaskOnly), as the shutdown strategy does) during the first exchange: only the first is removed;consumerProcessedStrategy=noneremoves nothing (control).Without the change 6 of the 7 fail (twice,
NeverWantedButInvokedfor the failed/undelivered rows, the document already removed while routing, and the removal exception thrown out ofpoll()); thenonecontrol passes.In-flight guard tests (same class; the couchbase route hands every exchange to
seda:async?waitForTaskToComplete=Never, and the route consuming that queue is started by the test when the handed-over exchanges should complete):maxMessagesPerPoll=1with three rows, and the consumer stopping mid batch: the undelivered rows are not left in flight, later polls deliver doc-2 and doc-3 while doc-1 is still in flight;consumerProcessedStrategy=nonestill delivers a document whose exchange is in flight (unchanged behaviour, control).Without the guard 6 of these 8 fail (the document delivered again while in flight); with the guard but without releasing the IDs, 4 fail (failure,
maxMessagesPerPoll, stopping, poll failing part way). With both commits the camel-couchbase module passes: 57 unit tests, 0 failures (the 8 Docker IT tests are skipped locally).Target
mainbranch)Tracking
Apache Camel coding standards and style
mvn clean install -DskipTestslocally from root folder and I have committed all auto-generated changes.(I built and tested the affected module, including the formatter and import-sort plugins. I did not run the full root build.)
AI-assisted contributions
Co-authored-bytrailers) and the PR description identifies the AI tool used.This PR was prepared with Claude Code (Claude Opus 5.5). Each commit carries a
Co-Authored-Bytrailer.Claude Code on behalf of allthingssecurity
🤖 Generated with Claude Code