davsclaus commented on code in PR #27644:
URL: https://github.com/apache/camel/pull/27644#discussion_r4236604583


##########
components/camel-jooq/src/test/java/org/apache/camel/component/jooq/JooqConsumerDeleteFailedTest.java:
##########
@@ -0,0 +1,98 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.jooq;
+
+import java.util.List;
+
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.ShutdownRunningTask;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.jooq.db.tables.records.AuthorRecord;
+import org.junit.jupiter.api.Test;
+
+import static org.apache.camel.component.jooq.db.Tables.AUTHOR;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * With consumeDelete (the default) an entity must only be deleted when its 
exchange was processed successfully: a
+ * failed or rollback only exchange, or an entity that was not processed, must 
leave the row in the table, so the next
+ * poll consumes it again.
+ */
+class JooqConsumerDeleteFailedTest extends BaseJooqTest {

Review Comment:
   The three cases look right and they fail without the fix. They all run with 
the default prototype exchange factory, so they would not catch the pooled 
exchange problem on `JooqConsumer` line 117. Could you add a variant (a 
subclass is enough) that overrides `createCamelContext()` and sets 
`getCamelContextExtension().setExchangeFactory(new PooledExchangeFactory())` 
before the routes are created? It fails with the current change and passes once 
the exchanges are created with `createExchange(false)`.



##########
components/camel-jooq/src/main/java/org/apache/camel/component/jooq/JooqConsumer.java:
##########
@@ -97,6 +114,7 @@ public int processBatch(Queue<Object> exchanges) throws 
Exception {
             for (int i = 0; i < total; i++) {
                 DataHolder holder = 
org.apache.camel.util.ObjectHelper.cast(DataHolder.class, exchanges.poll());
                 getProcessor().process(holder.exchange);
+                holder.consumed = !holder.exchange.isFailed() && 
!holder.exchange.isRollbackOnly();

Review Comment:
   This line runs after the unit of work is done. If the context uses the 
pooled exchange factory, `createExchange(true)` (line 99) makes 
`DefaultUnitOfWork.done()` call `PooledExchange.done()` at that point, and 
`done()` resets `exception` to null and `rollbackOnly` to false 
(`AbstractExchange.resetExtension()`). So here `isFailed()` and 
`isRollbackOnly()` are both false, and the failed entity is deleted anyway.
   
   camel-jpa (`JpaConsumer.processBatch`), camel-sql and camel-mybatis avoid 
this by creating the exchange without auto release and releasing it themselves 
after they have read the outcome:
   
   ```java
   // createExchange(Object result)
   Exchange exchange = createExchange(false);
   ```
   
   ```java
   DataHolder holder = 
org.apache.camel.util.ObjectHelper.cast(DataHolder.class, exchanges.poll());
   try {
       getProcessor().process(holder.exchange);
   } catch (Exception e) {
       holder.exchange.setException(e);
   }
   holder.consumed = !holder.exchange.isFailed() && 
!holder.exchange.isRollbackOnly();
   releaseExchange(holder.exchange, false);
   ```
   
   The try/catch (also from `JpaConsumer`) is optional but helps. Without it, 
an exception thrown straight out of `process()` stops the batch before the 
delete, so the entities that already succeeded are not deleted and get 
processed a second time on the next poll.



##########
components/camel-jooq/src/main/java/org/apache/camel/component/jooq/JooqConsumer.java:
##########
@@ -97,6 +114,7 @@ public int processBatch(Queue<Object> exchanges) throws 
Exception {
             for (int i = 0; i < total; i++) {

Review Comment:
   Optional: now that a record is deleted only when its exchange was processed, 
it is safe to check `isBatchAllowed()` on each iteration, as `JpaConsumer`, 
`SqlConsumer` and `MyBatisConsumer` do (`for (int index = 0; index < total && 
isBatchAllowed(); index++)`). Then, when a graceful shutdown starts in the 
middle of a batch, the remaining entities stay in the table for the next start 
instead of being processed while the route is stopping.



-- 
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]

Reply via email to