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]
