allthingssecurity commented on code in PR #27347:
URL: https://github.com/apache/camel/pull/27347#discussion_r4181400235
##########
components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java:
##########
@@ -277,6 +270,31 @@ private int pollWithView() throws Exception {
return processBatch(exchanges);
}
+ /**
+ * Removes the document once its exchange has been processed successfully.
+ * <p/>
+ * Removing it while the exchanges are built, before any of them is handed
to the route, loses the document when its
+ * exchange fails, or when it is never delivered at all because the batch
is cut short by {@code maxMessagesPerPoll}
+ * or by the consumer stopping. The on-completion runs only for an
exchange that went through the route and
+ * completed: a failed exchange keeps its document for the next poll, and
an undelivered one is not touched.
+ */
+ private void removeDocumentOnCompletion(Exchange exchange, String id) {
+ exchange.getExchangeExtension().addOnCompletion(new
SynchronizationAdapter() {
Review Comment:
Good point. I kept handover allowed: with `allowHandover()` returning
`false` the on-completion would run when the
consumer's own unit of work ends, so with `seda`
(`waitForTaskToComplete=Never`) or an async producer the document
would be removed before the asynchronous part had processed it, which is the
data loss this PR fixes.
On aws2-s3, which the PR description cites as the pattern: it does not
deliver the object twice. It also keeps
handover, but it adds the key to its `inProgressRepository` in the poll,
skips keys already there, and removes the key in
`onComplete`/`onFailure`. The file consumer does the same.
05b9c17f2451 adds the same guard here, for
`consumerProcessedStrategy=delete`: the consumer keeps the IDs of the
documents whose exchange is in flight and skips them in later polls. An ID
is released when its exchange completes or
fails, when it is never handed to the route (`maxMessagesPerPoll` or the
consumer stopping), and when the poll fails
part way. A skipped row does not count towards `maxMessagesPerPoll`. It is
an internal concurrent set, not a new
`inProgressRepository` option: it only has to cover this consumer's own
exchanges, and unlike aws2-s3's default
10000-entry FIFO cache it cannot evict an ID that is still in flight.
`none`/`filter` are unchanged, as they re-read the
document on every poll anyway.
Tests hand the exchange to `seda:...?waitForTaskToComplete=Never` and
complete it later: while it is in flight the next
poll does not deliver the document (SQL++ and view); after completion it was
processed and removed once; after failure
the next poll delivers it again; ids of undelivered rows and of a failed
poll are released. 6 fail without the guard,
4 fail without the release; module suite 57 tests, 0 failures.
Limit: the guard is per consumer instance. A consumer on another node
polling the same bucket can still read the
document before it is removed. The upgrade guide says so.
_Claude Code on behalf of allthingssecurity_
--
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]