gnodet-bot commented on code in PR #27314: URL: https://github.com/apache/camel/pull/27314#discussion_r4173397017
########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java: ########## @@ -0,0 +1,257 @@ +/* + * 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.hibernate; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Locale; +import java.util.Map; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.DefaultProducer; +import org.hibernate.Session; +import org.hibernate.SessionFactory; + +public class HibernateProducer extends DefaultProducer { + + private final HibernateEndpoint endpoint; + + public HibernateProducer(HibernateEndpoint endpoint) { + super(endpoint); + this.endpoint = endpoint; + } + + @Override + public void process(Exchange exchange) throws Exception { + if (endpoint.isJpaBacked()) { + processWithJpaTransaction(exchange); + } else { + processWithHibernateTransaction(exchange); + } + } + + private void processWithJpaTransaction(Exchange exchange) { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + exchange, + endpoint.getEntityManagerFactory(), + false, + true, + false); + + endpoint.getTransactionStrategy().executeInTransaction(() -> { + try { + doProcess(exchange, entityManager); + } catch (RuntimeException e) { + throw e; + } catch (Exception e) { + throw RuntimeCamelException.wrapRuntimeCamelException(e); + } + }); + } + + private void processWithHibernateTransaction(Exchange exchange) throws Exception { + SessionFactory sessionFactory = endpoint.getResolvedSessionFactory(); + + try (Session session = sessionFactory.openSession()) { + org.hibernate.Transaction transaction = session.beginTransaction(); + + try { + doProcess(exchange, session); + transaction.commit(); + } catch (Exception e) { + if (transaction.isActive()) { + transaction.rollback(); + } + throw e; + } + } + } + + private void doProcess(Exchange exchange, EntityManager entityManager) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, entityManager); + } else { + executeEntityOperation(exchange, entityManager); + } + } + + private void doProcess(Exchange exchange, Session session) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, session); + } else { + executeEntityOperation(exchange, session); + } + } + + private void executeEntityOperation(Exchange exchange, EntityManager entityManager) { + Object body = exchange.getIn().getBody(); + if (body == null) { + return; + } + + Object result; + if (body instanceof Collection<?> collection) { + List<Object> list = new ArrayList<>(collection.size()); + for (Object item : collection) { + if (endpoint.isUsePersist()) { + entityManager.persist(item); + list.add(item); + } else { + list.add(entityManager.merge(item)); + } + } + result = list; + } else { + if (endpoint.isUsePersist()) { + entityManager.persist(body); + result = body; + } else { + result = entityManager.merge(body); + } + } + + entityManager.flush(); + exchange.getMessage().setBody(result); + } + + private void executeEntityOperation(Exchange exchange, Session session) { + Object body = exchange.getIn().getBody(); + if (body == null) { + return; + } + + Object result; + if (body instanceof Collection<?> collection) { + List<Object> list = new ArrayList<>(collection.size()); + for (Object item : collection) { + if (endpoint.isUsePersist()) { + session.persist(item); + list.add(item); + } else { + list.add(session.merge(item)); + } + } + result = list; + } else { + if (endpoint.isUsePersist()) { + session.persist(body); + result = body; + } else { + result = session.merge(body); + } + } + + session.flush(); + exchange.getMessage().setBody(result); + } + + private void executeQuery(Exchange exchange, EntityManager entityManager) { + String hql = exchange.getIn().getHeader( + HibernateConstants.HIBERNATE_QUERY, + String.class); + + Query query; + + if (hql != null) { + query = entityManager.createQuery(hql); + } else if (endpoint.getNamedQuery() != null) { + query = entityManager.createNamedQuery(endpoint.getNamedQuery()); + } else if (endpoint.getNativeQuery() != null) { + query = entityManager.createNativeQuery(endpoint.getNativeQuery()); + } else { + query = entityManager.createQuery(endpoint.getQuery()); + } + + configureQuery(exchange, query); + + if (isExecuteUpdateQuery(hql, endpoint.getQuery())) { + int updated = query.executeUpdate(); + exchange.getMessage().setBody(updated); + } else { + exchange.getMessage().setBody(query.getResultList()); + } + } Review Comment: ⚠️ **Bug — `isExecuteUpdateQuery` doesn't detect mutation named/native queries:** When a `namedQuery` or `nativeQuery` is configured that performs DELETE/UPDATE, `isExecuteUpdateQuery(hql, endpoint.getQuery())` is called with `hql=null` (no header) and `endpoint.getQuery()=null` (the endpoint uses `namedQuery`/`nativeQuery`, not `query`). This returns `false`, so `query.getResultList()` is called on a mutation query — which will either fail or return unexpected results. The fix should also check `endpoint.getNamedQuery()` and `endpoint.getNativeQuery()` as fallbacks. However, for named queries, the query text isn't available at runtime — consider adding a `CamelHibernateQueryType` header or an endpoint option (e.g., `queryType=select|update`) to let users indicate mutation intent for named queries. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java: ########## @@ -0,0 +1,257 @@ +/* + * 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.hibernate; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Locale; +import java.util.Map; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.DefaultProducer; +import org.hibernate.Session; +import org.hibernate.SessionFactory; + +public class HibernateProducer extends DefaultProducer { + + private final HibernateEndpoint endpoint; + + public HibernateProducer(HibernateEndpoint endpoint) { + super(endpoint); + this.endpoint = endpoint; + } + + @Override + public void process(Exchange exchange) throws Exception { + if (endpoint.isJpaBacked()) { + processWithJpaTransaction(exchange); + } else { + processWithHibernateTransaction(exchange); + } + } + + private void processWithJpaTransaction(Exchange exchange) { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + exchange, + endpoint.getEntityManagerFactory(), + false, + true, + false); + + endpoint.getTransactionStrategy().executeInTransaction(() -> { + try { + doProcess(exchange, entityManager); + } catch (RuntimeException e) { + throw e; + } catch (Exception e) { + throw RuntimeCamelException.wrapRuntimeCamelException(e); + } + }); + } + + private void processWithHibernateTransaction(Exchange exchange) throws Exception { + SessionFactory sessionFactory = endpoint.getResolvedSessionFactory(); + + try (Session session = sessionFactory.openSession()) { + org.hibernate.Transaction transaction = session.beginTransaction(); + + try { + doProcess(exchange, session); + transaction.commit(); + } catch (Exception e) { + if (transaction.isActive()) { + transaction.rollback(); + } + throw e; + } + } + } + + private void doProcess(Exchange exchange, EntityManager entityManager) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, entityManager); + } else { + executeEntityOperation(exchange, entityManager); + } + } + + private void doProcess(Exchange exchange, Session session) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, session); + } else { + executeEntityOperation(exchange, session); + } + } + + private void executeEntityOperation(Exchange exchange, EntityManager entityManager) { + Object body = exchange.getIn().getBody(); + if (body == null) { + return; + } Review Comment: ⚠️ **Significant code duplication:** The Producer has near-identical method pairs: `doProcess` (×2), `executeEntityOperation` (×2), `executeQuery` (×2). Hibernate's `Session` extends JPA's `EntityManager` (since Hibernate 6+), and `org.hibernate.query.Query` extends `jakarta.persistence.Query`. You could eliminate the Session-specific methods entirely by using `session.unwrap(EntityManager.class)` or, more idiomatically, by always working through the `EntityManager` interface when an EMF is available and directly with the Session when not — but the Session *is* an EntityManager. This would cut ~100 lines of duplicated logic and eliminate the risk of the two paths diverging silently. Consider refactoring to a single code path. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateConsumer.java: ########## @@ -0,0 +1,393 @@ +/* + * 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.hibernate; + +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.util.Collection; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.Queue; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.ManyToOne; +import jakarta.persistence.OneToOne; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaConstants; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.ScheduledBatchPollingConsumer; +import org.apache.camel.util.CastUtils; +import org.hibernate.Session; +import org.hibernate.SessionFactory; +import org.hibernate.Transaction; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * The Hibernate consumer for polling database records via JPA or Native Hibernate Sessions. + */ +public class HibernateConsumer extends ScheduledBatchPollingConsumer { + + private static final Logger LOG = LoggerFactory.getLogger(HibernateConsumer.class); + + private final HibernateEndpoint endpoint; + + public HibernateConsumer(HibernateEndpoint endpoint, Processor processor) { + super(endpoint, processor); + this.endpoint = endpoint; + } + + @Override + protected int poll() throws Exception { + shutdownRunningTask = null; + pendingExchanges = 0; + + try { + if (endpoint.isJpaBacked()) { + return pollJpa(); + } else { + return pollNativeHibernate(); + } + } catch (Exception e) { + LOG.error("Polling failed on endpoint {}: {}", endpoint.getEndpointUri(), e.getMessage(), e); + throw e; + } + } + + private int pollJpa() throws Exception { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + null, + endpoint.getEntityManagerFactory(), + false, + true, + false); + + final int[] messagePolled = { 0 }; + + try { + endpoint.getTransactionStrategy().executeInTransaction(() -> { + try { + Queue<DataHolder> exchanges = createExchanges(entityManager); + messagePolled[0] = processBatch(CastUtils.cast(exchanges)); + + if (endpoint.isConsumeDelete()) { + entityManager.flush(); + } + } catch (RuntimeException e) { + throw e; + } catch (Exception e) { + throw RuntimeCamelException.wrapRuntimeCamelException(e); + } + }); + } catch (Exception e) { + LOG.error("JPA transaction processing failed: {}", e.getMessage(), e); + throw e; + } + + return messagePolled[0]; + } + + private int pollNativeHibernate() throws Exception { + SessionFactory sessionFactory = endpoint.getResolvedSessionFactory(); + + try (Session session = sessionFactory.openSession()) { + Transaction transaction = session.beginTransaction(); + + try { + Queue<DataHolder> exchanges = createExchanges(session); + int messagePolled = processBatch(CastUtils.cast(exchanges)); + + if (endpoint.isConsumeDelete()) { + session.flush(); + } + + transaction.commit(); + return messagePolled; + } catch (Exception e) { + if (transaction.isActive()) { + transaction.rollback(); + } + throw e; + } + } + } + + private Queue<DataHolder> createExchanges(EntityManager entityManager) { + Query query = createQuery(entityManager); + configureParameters(query); + + List<?> results = query.getResultList(); + forceConsumerAsReady(); + + Queue<DataHolder> exchanges = new LinkedList<>(); + for (Object result : results) { + Exchange exchange = createExchange(result, entityManager); + + DataHolder holder = new DataHolder(); + holder.exchange = exchange; + holder.entity = result; + holder.entityManager = entityManager; + + exchanges.add(holder); + } + + return exchanges; + } + + private Queue<DataHolder> createExchanges(Session session) { + org.hibernate.query.Query<?> query = createQuery(session); + configureParameters(query); + + List<?> results = query.getResultList(); + forceConsumerAsReady(); + + Queue<DataHolder> exchanges = new LinkedList<>(); + for (Object result : results) { + Exchange exchange = createExchange(result, null); + + DataHolder holder = new DataHolder(); + holder.exchange = exchange; + holder.entity = result; + holder.session = session; + + exchanges.add(holder); + } + + return exchanges; + } + + private Query createQuery(EntityManager entityManager) { + if (endpoint.getNamedQuery() != null) { + return entityManager.createNamedQuery(endpoint.getNamedQuery()); + } + + if (endpoint.getNativeQuery() != null) { + return entityManager.createNativeQuery(endpoint.getNativeQuery()); + } + + if (endpoint.getQuery() != null) { + return entityManager.createQuery(endpoint.getQuery()); + } + + if (endpoint.getEntityType() != null) { + return entityManager.createQuery("FROM " + endpoint.getEntityType().getName()); + } + + throw new IllegalArgumentException( + "No query or entityType configured for HibernateConsumer"); + } + + private org.hibernate.query.Query<?> createQuery(Session session) { + if (endpoint.getNamedQuery() != null) { + return session.createNamedQuery(endpoint.getNamedQuery(), Object.class); + } + + if (endpoint.getNativeQuery() != null) { + return session.createNativeQuery(endpoint.getNativeQuery(), Object.class); + } + + if (endpoint.getQuery() != null) { + return session.createQuery(endpoint.getQuery(), Object.class); + } + + if (endpoint.getEntityType() != null) { + return session.createQuery("FROM " + endpoint.getEntityType().getName(), Object.class); + } + + throw new IllegalArgumentException( + "No query or entityType configured for HibernateConsumer"); + } + + private void configureParameters(Query query) { + applyParametersToQuery(query::setParameter); + + if (endpoint.getMaximumResults() > 0) { + query.setMaxResults(endpoint.getMaximumResults()); + } + } + + private void configureParameters(org.hibernate.query.Query<?> query) { + applyParametersToQuery(query::setParameter); + + if (endpoint.getMaximumResults() > 0) { + query.setMaxResults(endpoint.getMaximumResults()); + } + } + + @FunctionalInterface + private interface ParameterBinder { + void bind(String name, Object value); + } + + private void applyParametersToQuery(ParameterBinder binder) { + Map<String, Object> params = endpoint.getParameters(); + if (params != null) { + for (Map.Entry<String, Object> entry : params.entrySet()) { + binder.bind(entry.getKey(), entry.getValue()); + } + } + } + + @Override + public int processBatch(Queue<Object> exchanges) throws Exception { + int total = exchanges.size(); + + for (Object exchangeObject : exchanges) { + DataHolder holder = (DataHolder) exchangeObject; + Exchange exchange = holder.exchange; + + try { + getProcessor().process(exchange); + + if (exchange.getException() != null) { + throw exchange.getException(); + } + + if (endpoint.isConsumeDelete()) { + deleteEntity(holder); + } + } finally { + releaseExchange(exchange, false); + } + } + + return total; + } + + private void deleteEntity(DataHolder holder) { + if (holder.entity == null) { + return; + } + + if (holder.entityManager != null) { + EntityManager entityManager = holder.entityManager; + Object entityToDelete = holder.entity; + + if (!entityManager.contains(entityToDelete)) { + entityToDelete = entityManager.merge(entityToDelete); + } + + detachRelationshipsGenerically(entityToDelete); + entityManager.remove(entityToDelete); + + } else if (holder.session != null) { + Session session = holder.session; + Object entityToDelete = holder.entity; + + if (!session.contains(entityToDelete)) { + entityToDelete = session.merge(entityToDelete); + } + + detachRelationshipsGenerically(entityToDelete); + session.remove(entityToDelete); + } + } + + /** + * Inspects entity fields dynamically via reflection for standard JPA relationship annotations + * (@ManyToOne, @OneToOne) and detaches the entity from parent collections to prevent re-persisting during flush. + */ + private void detachRelationshipsGenerically(Object entity) { + if (entity == null) { + return; + } + + Class<?> clazz = entity.getClass(); + while (clazz != null && clazz != Object.class) { + for (Field field : clazz.getDeclaredFields()) { + if (field.isAnnotationPresent(ManyToOne.class) || field.isAnnotationPresent(OneToOne.class)) { + try { + field.setAccessible(true); + Object parent = field.get(entity); + if (parent != null) { + removeFromParentCollections(parent, entity); + field.set(entity, null); + } + } catch (Exception e) { + LOG.debug("Could not automatically unbind relationship field '{}' on entity {}: {}", + field.getName(), entity.getClass().getName(), e.getMessage()); + } + } + } + clazz = clazz.getSuperclass(); Review Comment: ⚠️ **`detachRelationshipsGenerically` is fragile and incomplete:** 1. **Missing `@ManyToMany`:** Only `@ManyToOne` and `@OneToOne` are handled. `@ManyToMany` bidirectional relationships have the same re-persistence risk. 2. **Hibernate proxies:** `field.get(entity)` on a lazy `@ManyToOne` will trigger proxy initialization and a potential `LazyInitializationException` if the session is in a bad state. Consider checking `Hibernate.isInitialized()` before accessing. 3. **Security manager:** `field.setAccessible(true)` may fail under strict security policies. The catch block logs at DEBUG and continues — this is acceptable, but the overall approach is inherently brittle. 4. **`removeFromParentCollections` walks getters via `getMethods()`:** This includes inherited methods from Object and potentially unrelated getters returning collections. Filter to methods declared on the actual entity class or annotated with `@OneToMany`/`@ManyToMany`. 5. **Swallowed exceptions in `removeFromParentCollections`:** `catch (Exception ignored)` is too broad — a `ConcurrentModificationException` from modifying a collection during iteration would be silently lost. Consider making the relationship-detach opt-in rather than a default behavior, or at least document this reflection-based approach clearly so users understand the limitations. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java: ########## @@ -0,0 +1,257 @@ +/* + * 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.hibernate; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Locale; +import java.util.Map; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.DefaultProducer; +import org.hibernate.Session; +import org.hibernate.SessionFactory; + +public class HibernateProducer extends DefaultProducer { + + private final HibernateEndpoint endpoint; + + public HibernateProducer(HibernateEndpoint endpoint) { + super(endpoint); + this.endpoint = endpoint; + } + + @Override + public void process(Exchange exchange) throws Exception { + if (endpoint.isJpaBacked()) { + processWithJpaTransaction(exchange); + } else { + processWithHibernateTransaction(exchange); + } + } + + private void processWithJpaTransaction(Exchange exchange) { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + exchange, + endpoint.getEntityManagerFactory(), + false, Review Comment: 💡 **Potential EntityManager leak in JPA path:** `JpaHelper.getTargetEntityManager()` creates an EntityManager but `processWithJpaTransaction` never explicitly closes it. In the `DefaultTransactionStrategy` flow, the EntityManager may be managed by Spring's transaction synchronization — but if no Spring transaction infrastructure is present (pure JPA, no Spring), this EntityManager will leak. Compare with `camel-jpa`'s `JpaProducer`, which calls `entityManager.close()` in a `finally` block when `usePassedInEntityManager` is false. Ensure the same cleanup contract applies here. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateComponent.java: ########## @@ -0,0 +1,128 @@ +/* + * 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.hibernate; + +import java.util.Map; + +import jakarta.persistence.EntityManagerFactory; + +import org.apache.camel.Endpoint; +import org.apache.camel.spi.Metadata; +import org.apache.camel.spi.annotations.Component; +import org.apache.camel.support.DefaultComponent; +import org.hibernate.SessionFactory; +import org.springframework.transaction.PlatformTransactionManager; + +/** + * Native Hibernate component for Apache Camel supporting both JPA-backed and Native Hibernate execution modes. + */ +@Component("hibernate") +public class HibernateComponent extends DefaultComponent { + + @Metadata(description = "To use the SessionFactory as the factory for creating Hibernate sessions.") + private SessionFactory sessionFactory; + Review Comment: 💡 **Duplicate `@Metadata` annotations:** The `@Metadata` annotation with the description is declared on both the fields (lines 30-37) and the setters (lines 94, 106, 118). Camel's code generator uses the field-level annotation — the setter-level annotations are redundant and will get out of sync with refactoring. ########## components/camel-hibernate/pom.xml: ########## @@ -0,0 +1,91 @@ +<?xml version="1.0" encoding="UTF-8"?> +<!-- + + 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. + +--> +<project xmlns="http://maven.apache.org/POM/4.0.0" Review Comment: 🔴 **Critical — `firstVersion` mismatch:** The POM declares `<firstVersion>4.11.0</firstVersion>` but the `@UriEndpoint` annotation on `HibernateEndpoint.java` declares `firstVersion = "4.23.0"`. Since this is a brand new component targeting 4.23.0, the POM property is wrong and will cause incorrect metadata in the Camel catalog. ```suggestion <firstVersion>4.23.0</firstVersion> ``` ########## components/camel-hibernate/pom.xml: ########## @@ -0,0 +1,91 @@ +<?xml version="1.0" encoding="UTF-8"?> +<!-- + + 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. + +--> +<project xmlns="http://maven.apache.org/POM/4.0.0" + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> + <modelVersion>4.0.0</modelVersion> + + <parent> + <groupId>org.apache.camel</groupId> + <artifactId>components</artifactId> + <version>4.23.0-SNAPSHOT</version> + </parent> + + <artifactId>camel-hibernate</artifactId> + <packaging>jar</packaging> + <name>Camel :: Hibernate</name> + <description>Camel Hibernate Component</description> + + <properties> + <firstVersion>4.11.0</firstVersion> + <label>database</label> + </properties> + + <dependencies> + <!-- Camel Support --> Review Comment: 💡 **Explicit version overrides:** The dependencies for `jakarta.persistence-api`, `hibernate-core`, `spring-orm`, and `h2` all use explicit `${...}` version properties. Camel's parent POM already manages these — using explicit versions creates a risk of version drift. Use the managed versions instead: ```suggestion <dependency> <groupId>jakarta.persistence</groupId> <artifactId>jakarta.persistence-api</artifactId> </dependency> ``` Same applies to `hibernate-core`, `spring-orm`, and `h2` — drop the explicit `<version>` tags. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateEndpoint.java: ########## @@ -0,0 +1,293 @@ +/* + * 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.hibernate; + +import java.util.Map; + +import jakarta.persistence.EntityManagerFactory; + +import org.apache.camel.Category; +import org.apache.camel.Consumer; +import org.apache.camel.Processor; +import org.apache.camel.Producer; +import org.apache.camel.component.jpa.DefaultTransactionStrategy; +import org.apache.camel.component.jpa.TransactionStrategy; +import org.apache.camel.spi.Metadata; +import org.apache.camel.spi.UriEndpoint; +import org.apache.camel.spi.UriParam; +import org.apache.camel.spi.UriPath; +import org.apache.camel.support.ScheduledPollEndpoint; +import org.hibernate.SessionFactory; +import org.springframework.transaction.PlatformTransactionManager; + +/** + * Perform database operations using Hibernate ORM supporting both JPA-backed and Native Hibernate modes. + */ +@UriEndpoint(firstVersion = "4.23.0", scheme = "hibernate", title = "Hibernate", syntax = "hibernate:entityClassName", + category = { Category.DATABASE }) +public class HibernateEndpoint extends ScheduledPollEndpoint { + + @UriPath(description = "Target entity class name or entity type name") + @Metadata(required = true) + private String entityClassName; + + private Class<?> entityType; + + @UriParam(description = "The EntityManagerFactory to use") + private EntityManagerFactory entityManagerFactory; + + @UriParam(description = "The Hibernate SessionFactory to use") + private SessionFactory sessionFactory; + + @UriParam(description = "The PlatformTransactionManager to use") + private PlatformTransactionManager transactionManager; + + private TransactionStrategy transactionStrategy; + + @UriParam(description = "HQL query to execute") + private String query; + + @UriParam(description = "Named query to execute") + private String namedQuery; + + @UriParam(description = "Native SQL query to execute") + private String nativeQuery; + + @UriParam(defaultValue = "-1", description = "Maximum number of results to retrieve") + private int maximumResults = -1; + + @UriParam(defaultValue = "true", description = "Whether to delete consumed entities after polling") + private boolean consumeDelete = true; + + @UriParam(defaultValue = "false", + description = "Indicates to use entityManager.persist(entity) or session.persist(entity) instead of merge") + private boolean usePersist; + + @UriParam(description = "Parameters to pass to the query in key-value map format", multiValue = true, + prefix = "parameters.") + private Map<String, Object> parameters; + + public HibernateEndpoint() { + } + + public HibernateEndpoint(String uri, HibernateComponent component) { + super(uri, component); + } + + @Override + public Producer createProducer() throws Exception { + return new HibernateProducer(this); + } + + @Override + public Consumer createConsumer(Processor processor) throws Exception { + HibernateConsumer consumer = new HibernateConsumer(this, processor); + configureConsumer(consumer); + return consumer; + } + + public SessionFactory getResolvedSessionFactory() { + if (sessionFactory != null) { + return sessionFactory; + } + + if (entityManagerFactory != null) { + return entityManagerFactory.unwrap(SessionFactory.class); + } + + throw new IllegalArgumentException( + "Either SessionFactory or EntityManagerFactory must be configured on HibernateEndpoint"); + } + + public boolean isJpaBacked() { + return entityManagerFactory != null; + } + + public boolean isNativeHibernate() { + return entityManagerFactory == null && sessionFactory != null; + } + + public TransactionStrategy getTransactionStrategy() { + if (!isJpaBacked()) { + throw new IllegalStateException( + "A TransactionStrategy is only available when an EntityManagerFactory is configured"); + } + + if (transactionStrategy == null) { + transactionStrategy = createTransactionStrategy(); Review Comment: 💡 **Race condition in lazy `TransactionStrategy` init:** `getTransactionStrategy()` creates the strategy lazily without synchronization. If two threads call `createProducer()` / `createConsumer()` concurrently during endpoint startup, they could both see `transactionStrategy == null` and create separate instances. In practice, `DefaultTransactionStrategy` is stateless enough that this is unlikely to cause bugs, but it violates the principle of least surprise. Consider using a `volatile` field + double-checked locking, or initialize it eagerly in `doStart()`. ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateConsumer.java: ########## @@ -0,0 +1,393 @@ +/* + * 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.hibernate; + +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.util.Collection; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.Queue; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.ManyToOne; +import jakarta.persistence.OneToOne; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaConstants; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.ScheduledBatchPollingConsumer; +import org.apache.camel.util.CastUtils; +import org.hibernate.Session; +import org.hibernate.SessionFactory; +import org.hibernate.Transaction; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * The Hibernate consumer for polling database records via JPA or Native Hibernate Sessions. + */ +public class HibernateConsumer extends ScheduledBatchPollingConsumer { + + private static final Logger LOG = LoggerFactory.getLogger(HibernateConsumer.class); + + private final HibernateEndpoint endpoint; + + public HibernateConsumer(HibernateEndpoint endpoint, Processor processor) { + super(endpoint, processor); + this.endpoint = endpoint; + } + + @Override + protected int poll() throws Exception { + shutdownRunningTask = null; + pendingExchanges = 0; + + try { + if (endpoint.isJpaBacked()) { + return pollJpa(); + } else { + return pollNativeHibernate(); + } + } catch (Exception e) { + LOG.error("Polling failed on endpoint {}: {}", endpoint.getEndpointUri(), e.getMessage(), e); + throw e; + } + } + + private int pollJpa() throws Exception { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + null, + endpoint.getEntityManagerFactory(), + false, + true, + false); + + final int[] messagePolled = { 0 }; + + try { + endpoint.getTransactionStrategy().executeInTransaction(() -> { + try { + Queue<DataHolder> exchanges = createExchanges(entityManager); + messagePolled[0] = processBatch(CastUtils.cast(exchanges)); + + if (endpoint.isConsumeDelete()) { + entityManager.flush(); + } + } catch (RuntimeException e) { + throw e; + } catch (Exception e) { + throw RuntimeCamelException.wrapRuntimeCamelException(e); + } + }); + } catch (Exception e) { + LOG.error("JPA transaction processing failed: {}", e.getMessage(), e); + throw e; + } + + return messagePolled[0]; + } + + private int pollNativeHibernate() throws Exception { + SessionFactory sessionFactory = endpoint.getResolvedSessionFactory(); + + try (Session session = sessionFactory.openSession()) { + Transaction transaction = session.beginTransaction(); + + try { + Queue<DataHolder> exchanges = createExchanges(session); + int messagePolled = processBatch(CastUtils.cast(exchanges)); + + if (endpoint.isConsumeDelete()) { + session.flush(); + } + + transaction.commit(); + return messagePolled; + } catch (Exception e) { + if (transaction.isActive()) { + transaction.rollback(); + } + throw e; + } + } + } + + private Queue<DataHolder> createExchanges(EntityManager entityManager) { + Query query = createQuery(entityManager); + configureParameters(query); + + List<?> results = query.getResultList(); + forceConsumerAsReady(); + + Queue<DataHolder> exchanges = new LinkedList<>(); + for (Object result : results) { + Exchange exchange = createExchange(result, entityManager); + + DataHolder holder = new DataHolder(); + holder.exchange = exchange; + holder.entity = result; + holder.entityManager = entityManager; + + exchanges.add(holder); + } + + return exchanges; + } + + private Queue<DataHolder> createExchanges(Session session) { + org.hibernate.query.Query<?> query = createQuery(session); + configureParameters(query); + + List<?> results = query.getResultList(); + forceConsumerAsReady(); + + Queue<DataHolder> exchanges = new LinkedList<>(); + for (Object result : results) { + Exchange exchange = createExchange(result, null); + + DataHolder holder = new DataHolder(); + holder.exchange = exchange; + holder.entity = result; + holder.session = session; + + exchanges.add(holder); + } + + return exchanges; + } + + private Query createQuery(EntityManager entityManager) { + if (endpoint.getNamedQuery() != null) { + return entityManager.createNamedQuery(endpoint.getNamedQuery()); + } + + if (endpoint.getNativeQuery() != null) { + return entityManager.createNativeQuery(endpoint.getNativeQuery()); + } + + if (endpoint.getQuery() != null) { + return entityManager.createQuery(endpoint.getQuery()); + } + + if (endpoint.getEntityType() != null) { + return entityManager.createQuery("FROM " + endpoint.getEntityType().getName()); + } + + throw new IllegalArgumentException( + "No query or entityType configured for HibernateConsumer"); + } + + private org.hibernate.query.Query<?> createQuery(Session session) { + if (endpoint.getNamedQuery() != null) { + return session.createNamedQuery(endpoint.getNamedQuery(), Object.class); + } + + if (endpoint.getNativeQuery() != null) { + return session.createNativeQuery(endpoint.getNativeQuery(), Object.class); + } + + if (endpoint.getQuery() != null) { + return session.createQuery(endpoint.getQuery(), Object.class); + } + + if (endpoint.getEntityType() != null) { + return session.createQuery("FROM " + endpoint.getEntityType().getName(), Object.class); + } + + throw new IllegalArgumentException( + "No query or entityType configured for HibernateConsumer"); + } + + private void configureParameters(Query query) { + applyParametersToQuery(query::setParameter); + + if (endpoint.getMaximumResults() > 0) { + query.setMaxResults(endpoint.getMaximumResults()); + } + } + + private void configureParameters(org.hibernate.query.Query<?> query) { + applyParametersToQuery(query::setParameter); + + if (endpoint.getMaximumResults() > 0) { + query.setMaxResults(endpoint.getMaximumResults()); + } + } + + @FunctionalInterface + private interface ParameterBinder { + void bind(String name, Object value); + } + + private void applyParametersToQuery(ParameterBinder binder) { + Map<String, Object> params = endpoint.getParameters(); + if (params != null) { + for (Map.Entry<String, Object> entry : params.entrySet()) { + binder.bind(entry.getKey(), entry.getValue()); + } + } + } + + @Override + public int processBatch(Queue<Object> exchanges) throws Exception { + int total = exchanges.size(); + + for (Object exchangeObject : exchanges) { + DataHolder holder = (DataHolder) exchangeObject; + Exchange exchange = holder.exchange; + + try { + getProcessor().process(exchange); + + if (exchange.getException() != null) { + throw exchange.getException(); + } + + if (endpoint.isConsumeDelete()) { + deleteEntity(holder); + } + } finally { + releaseExchange(exchange, false); + } + } Review Comment: ⚠️ **Missing `isBatchAllowed()` check in `processBatch`:** Other Camel polling consumers (e.g., `JpaConsumer`, `FileConsumer`) check `isBatchAllowed()` before processing each exchange to support graceful shutdown. Without this check, the consumer will continue processing all queued entities even when the route is being stopped, which can cause shutdown delays or data inconsistency if the transaction is rolled back mid-batch. ```suggestion public int processBatch(Queue<Object> exchanges) throws Exception { int total = exchanges.size(); int processed = 0; for (Object exchangeObject : exchanges) { if (!isBatchAllowed()) { break; } DataHolder holder = (DataHolder) exchangeObject; Exchange exchange = holder.exchange; try { getProcessor().process(exchange); if (exchange.getException() != null) { throw exchange.getException(); } if (endpoint.isConsumeDelete()) { deleteEntity(holder); } } finally { releaseExchange(exchange, false); } processed++; } return processed; } ``` ########## components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java: ########## @@ -0,0 +1,257 @@ +/* + * 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.hibernate; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Locale; +import java.util.Map; + +import jakarta.persistence.EntityManager; +import jakarta.persistence.Query; + +import org.apache.camel.Exchange; +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.component.jpa.JpaHelper; +import org.apache.camel.support.DefaultProducer; +import org.hibernate.Session; +import org.hibernate.SessionFactory; + +public class HibernateProducer extends DefaultProducer { + + private final HibernateEndpoint endpoint; + + public HibernateProducer(HibernateEndpoint endpoint) { + super(endpoint); + this.endpoint = endpoint; + } + + @Override + public void process(Exchange exchange) throws Exception { + if (endpoint.isJpaBacked()) { + processWithJpaTransaction(exchange); + } else { + processWithHibernateTransaction(exchange); + } + } + + private void processWithJpaTransaction(Exchange exchange) { + EntityManager entityManager = JpaHelper.getTargetEntityManager( + exchange, + endpoint.getEntityManagerFactory(), + false, + true, + false); + + endpoint.getTransactionStrategy().executeInTransaction(() -> { + try { + doProcess(exchange, entityManager); + } catch (RuntimeException e) { + throw e; + } catch (Exception e) { + throw RuntimeCamelException.wrapRuntimeCamelException(e); + } + }); + } + + private void processWithHibernateTransaction(Exchange exchange) throws Exception { + SessionFactory sessionFactory = endpoint.getResolvedSessionFactory(); + + try (Session session = sessionFactory.openSession()) { + org.hibernate.Transaction transaction = session.beginTransaction(); + + try { + doProcess(exchange, session); + transaction.commit(); + } catch (Exception e) { + if (transaction.isActive()) { + transaction.rollback(); + } + throw e; + } + } + } + + private void doProcess(Exchange exchange, EntityManager entityManager) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, entityManager); + } else { + executeEntityOperation(exchange, entityManager); + } + } + + private void doProcess(Exchange exchange, Session session) { + if (endpoint.getQuery() != null + || endpoint.getNamedQuery() != null + || endpoint.getNativeQuery() != null + || exchange.getIn().getHeader(HibernateConstants.HIBERNATE_QUERY) != null) { + executeQuery(exchange, session); + } else { + executeEntityOperation(exchange, session); Review Comment: 💡 **Silent no-op on null body:** When `body == null`, `executeEntityOperation` returns silently without setting any body on the outgoing message. This means the exchange passes through with whatever body was previously set, which could confuse downstream processors expecting a result. Consider at minimum logging a warning, or setting the body to null explicitly on the out message to make the behavior deterministic. ########## components/camel-hibernate/src/main/docs/hibernate-component.adoc: ########## @@ -0,0 +1,212 @@ += Hibernate Component Review Comment: ⚠️ **Broken AsciiDoc formatting:** The documentation uses Markdown `#` headers instead of AsciiDoc `=` headers throughout. The tab blocks (`[tabs]`) are malformed — they use `##` instead of `====` for tab delimiters. The Maven dependency XML snippet and several code blocks are empty or incomplete. The XML and YAML examples in multiple sections are blank. This will render incorrectly in the Camel website. Please regenerate or manually fix the AsciiDoc formatting to match other Camel component docs (e.g., `camel-jpa/src/main/docs/jpa-component.adoc`). -- 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]
