gnodet-bot commented on code in PR #27314:
URL: https://github.com/apache/camel/pull/27314#discussion_r4181835773
##########
parent/pom.xml:
##########
@@ -239,6 +240,7 @@
<hdrhistrogram-version>2.2.2</hdrhistrogram-version>
<hibernate-validator-version>9.1.4.Final</hibernate-validator-version>
<hibernate-version>7.4.11.Final</hibernate-version>
+ <hibernate8-version>8.0.0.Beta3</hibernate8-version>
Review Comment:
🔴 **Critical — Hibernate 8.0.0.Beta3 is pre-release.** This is a Beta that
is not yet GA. Shipping a Camel component on a Beta dependency means:
- API breakage on the next Hibernate 8 release will break camel-hibernate
- Users inheriting from the parent POM now see `hibernate8-version` and
`jakarta-persistence4-api-version` properties even though no other module uses
them
- JPA 4.0.0-M7 (line 189) is a milestone, not a release
If Hibernate 8 is the target, this should be documented as
experimental/preview in the component description and the JIRA. Consider
keeping these version properties in the component POM rather than polluting the
shared parent.
##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateConsumer.java:
##########
@@ -0,0 +1,104 @@
+/*
+ * 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.List;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.support.ScheduledPollConsumer;
+import org.hibernate.LockMode;
+import org.hibernate.Session;
+import org.hibernate.Timeouts;
+import org.hibernate.Transaction;
+import org.hibernate.query.SelectionQuery;
+
+public class HibernateConsumer extends ScheduledPollConsumer {
+
+ private final HibernateEndpoint endpoint;
+
+ public HibernateConsumer(HibernateEndpoint endpoint, Processor processor) {
+ super(endpoint, processor);
+ this.endpoint = endpoint;
+ }
+
+ @Override
+ protected int poll() throws Exception {
+ if (endpoint.getSelectionQuery() == null ||
endpoint.getSelectionQuery().isBlank()) {
+ throw new IllegalArgumentException("Hibernate consumer requires
selectionQuery");
+ }
+
+ Session session = endpoint.getTenantIdentifier() == null
+ ? endpoint.getSessionFactory().openSession()
+ : endpoint.getSessionFactory().withOptions()
+ .tenantIdentifier(endpoint.getTenantIdentifier())
+ .openSession();
+
+ try (session) {
+ Transaction transaction = session.beginTransaction();
+
+ try {
+ if (endpoint.getFilters() != null) {
+ endpoint.getFilters().forEach((filterName, parameters) -> {
+ var filter = session.enableFilter(filterName);
+ if (parameters != null) {
+ parameters.forEach(filter::setParameter);
+ }
+ });
+ }
+
+ SelectionQuery<?> query = session.createSelectionQuery(
+ endpoint.getSelectionQuery(),
endpoint.getEntityType());
+
+ if (endpoint.isReadOnly()) {
+ session.setDefaultReadOnly(true);
+ query.setReadOnly(true);
+ }
+
+ if (endpoint.getMaximumResults() > 0) {
+ query.setMaxResults(endpoint.getMaximumResults());
+ }
+
+ if (endpoint.isSkipLocked()) {
+ query.setHibernateLockMode(LockMode.PESSIMISTIC_WRITE);
+ query.setLockTimeout(Timeouts.SKIP_LOCKED);
+ }
+
+ List<?> results = query.getResultList();
+
+ for (Object result : results) {
+ Exchange exchange = createExchange(false);
+ exchange.getMessage().setBody(result);
+
+ getProcessor().process(exchange);
+
+ if (exchange.getException() != null) {
+ throw exchange.getException();
+ }
+ }
+
Review Comment:
⚠️ **Consumer exchange leak and poison-row problem.**
Two issues:
1. **Exchange leak:** `createExchange(false)` borrows an exchange from the
pool but `releaseExchange()` is never called. After processing, the exchange
must be released in a `finally` block.
2. **Poison row blocks the consumer:** If one entity causes
`getProcessor().process(exchange)` to fail (via `exchange.getException()`), the
exception is thrown, the entire transaction is rolled back, and **all**
entities — including those already successfully processed — will be re-polled.
A single bad entity permanently blocks the consumer.
The standard Camel pattern (see `JpaConsumer`, `FileConsumer`) processes
each exchange independently with its own error handling:
```suggestion
for (Object result : results) {
Exchange exchange = createExchange(false);
exchange.getMessage().setBody(result);
try {
getProcessor().process(exchange);
} catch (Exception e) {
handleException("Error processing exchange",
exchange, e);
} finally {
releaseExchange(exchange, false);
}
}
```
##########
components/camel-hibernate/src/main/resources/META-INF/services/org/apache/camel/component/hibernate:
##########
@@ -0,0 +1 @@
+class=org.apache.camel.component.hibernate.HibernateComponent
Review Comment:
⚠️ **Duplicate of generated file.** This hand-written service file
(`src/main/resources/...`) duplicates the generated one
(`src/generated/resources/...`) and has no license header. The generated file
already contains the same content with a `# Generated by camel build tools`
comment. Please remove this file — it was flagged in davsclaus's finding #9 and
remains unaddressed.
##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java:
##########
@@ -0,0 +1,159 @@
+/*
+ * 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 java.util.stream.Stream;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.support.DefaultProducer;
+import org.hibernate.KeyType;
+import org.hibernate.Session;
+import org.hibernate.StatelessSession;
+import org.hibernate.Transaction;
+import org.hibernate.query.MutationQuery;
+import org.hibernate.query.SelectionQuery;
+
+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.getStatelessOperation() != null) {
+ processStateless(exchange);
+ return;
+ }
+
+ Session session = endpoint.getTenantIdentifier() == null
+ ? endpoint.getSessionFactory().openSession()
+ : endpoint.getSessionFactory().withOptions()
+ .tenantIdentifier(endpoint.getTenantIdentifier())
+ .openSession();
+
+ Transaction transaction = session.beginTransaction();
+
+ try {
+ if (endpoint.getFilters() != null) {
+ endpoint.getFilters().forEach((filterName, parameters) -> {
+ var filter = session.enableFilter(filterName);
+ if (parameters != null) {
+ parameters.forEach(filter::setParameter);
+ }
+ });
+ }
+
+ @SuppressWarnings("unchecked")
+ Map<String, Object> parameters = exchange.getMessage().getHeader(
+ HibernateConstants.HIBERNATE_PARAMETERS, Map.class);
+
+ if (endpoint.getNaturalIdParameters() != null) {
+ Object entity = session.find(
+ endpoint.getEntityType(),
+ endpoint.getNaturalIdParameters(),
+ KeyType.NATURAL);
+
+ exchange.getMessage().setBody(entity);
+
+ transaction.commit();
+ session.close();
+ } else if (endpoint.getSelectionQuery() != null) {
+ SelectionQuery<?> query = session.createSelectionQuery(
+ endpoint.getSelectionQuery(),
endpoint.getEntityType());
+
+ session.setDefaultReadOnly(endpoint.isReadOnly());
+ query.setReadOnly(endpoint.isReadOnly());
+
+ if (parameters != null) {
+ parameters.forEach(query::setParameter);
+ }
+
+ if (endpoint.isStreaming()) {
+ Stream<?> stream = query.getResultStream();
+
+ exchange.getMessage().setBody(stream.onClose(
+ () -> closeStreamingSession(session,
transaction)));
+
+ return;
Review Comment:
⚠️ **Streaming session/transaction leak if the consumer never closes the
stream.** The session cleanup relies entirely on `stream.onClose()`. If the
downstream route does not close the `Stream` (which is easy to forget —
collecting to a list via `stream.toList()` does not trigger `onClose`), the
session and transaction remain open indefinitely.
Consider at minimum documenting this contract clearly (the caller MUST close
the stream), or adding a timeout/safety-net mechanism. Also, `stream.onClose()`
returns a *new* stream — verify this is the one that gets set as the body, not
the original (it is correct here, but worth a comment for maintainability).
##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateEndpoint.java:
##########
@@ -0,0 +1,250 @@
+/*
+ * 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 org.apache.camel.Category;
+import org.apache.camel.Consumer;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+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;
+
+@UriEndpoint(firstVersion = "4.23.0", scheme = "hibernate", title =
"Hibernate", syntax = "hibernate:entityClassName",
+ category = { Category.DATABASE })
Review Comment:
💡 **Missing `headersClass`.** The `@UriEndpoint` annotation should declare
`headersClass = HibernateConstants.class` so the catalog knows about the
`CamelHibernateParameters` header. Also add `@Metadata(label = "producer",
description = "...")` on the constant in `HibernateConstants`. This was
davsclaus's finding #8.
##########
components/camel-hibernate/src/main/docs/hibernate-component.adoc:
##########
@@ -0,0 +1,745 @@
+= Hibernate Component
+:doctitle: Hibernate
+:shortname: hibernate
+:artifactid: camel-hibernate
+:description: Camel Hibernate Component
+:since: 4.23
+:supportlevel: Preview
+:tabs-sync-option:
+:component-header: Both producer and consumer are supported
+
+*Since Camel {since}*
+
+*{component-header}*
+
+The Hibernate component provides integration with relational databases using
**Hibernate ORM 8**.
+
+The component uses Hibernate's native `SessionFactory`, `Session`,
`StatelessSession`, `SelectionQuery`, and `MutationQuery` APIs.
+
+It supports HQL selection and mutation queries, natural-id lookups, Hibernate
filters, read-only queries, streaming query results, stateless insert and
upsert operations, multi-tenancy, and competing consumers using `SKIP_LOCKED`.
+
+The component can either reuse an existing Hibernate `SessionFactory` or
bootstrap one using a configured `DataSource` and explicit entity classes.
+
+The component does not use Jakarta Persistence (JPA), `camel-jpa`, or Spring
`PlatformTransactionManager`.
+
+Maven users will need to add the following dependency to their `pom.xml`
+for this component:
+
+## [source,xml]
Review Comment:
⚠️ **Broken AsciiDoc formatting (still).** The documentation uses `##`
(Markdown-style headers) instead of AsciiDoc `====` tab delimiters and `----`
code fences throughout. For example:
- `## [source,xml]` should be `[source,xml]` followed by `----`
- `# [NOTE]` should be `[NOTE]` followed by `====`
- `## context.addComponent(...)` is raw code that should be inside a code
block
This renders incorrectly on the Camel website. Please follow the formatting
of other component docs (e.g., `camel-jpa/src/main/docs/jpa-component.adoc`).
This was davsclaus's finding #9 (docs), and the rewrite made the same
Markdown-style formatting errors.
##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateEndpoint.java:
##########
@@ -0,0 +1,250 @@
+/*
+ * 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 org.apache.camel.Category;
+import org.apache.camel.Consumer;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+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;
+
+@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;
+
+ @UriParam(description = "The HQL selection query to execute.")
+ private String selectionQuery;
+
+ @UriParam(description = "The HQL mutation query to execute.")
+ private String mutationQuery;
+
+ @UriParam(description = "The natural-id property values used for lookup.")
+ private Map<String, Object> naturalIdParameters;
+
+ @UriParam(description = "Whether the Hibernate session and selection query
should be read-only.")
+ private boolean readOnly;
+
+ @UriParam(description = "Hibernate filters and their parameter values.")
+ private Map<String, Map<String, Object>> filters;
+
+ @UriParam(description = "The tenant identifier used to create the
Hibernate session.")
+ private String tenantIdentifier;
+
+ @UriParam(description = "Stateless operation to perform: insert or
upsert.")
+ private String statelessOperation;
+
+ @UriParam(description = "Whether selection query results should be
returned as a stream.")
+ private boolean streaming;
+
+ @UriParam(description = "Whether the consumer should skip rows that are
already locked by another consumer.")
+ private boolean skipLocked;
+
+ @UriParam(description = "The maximum number of entities to retrieve in a
single poll.")
+ private int maximumResults;
+
+ private Class<?> entityType;
+ private SessionFactory sessionFactory;
+
+ 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;
+ }
+
+ @Override
+ protected void doStart() throws Exception {
+ if (sessionFactory == null) {
+ sessionFactory = ((HibernateComponent)
getComponent()).getSessionFactory();
+ }
+ if (sessionFactory == null) {
+ throw new IllegalArgumentException("SessionFactory must be
configured or available on HibernateComponent");
+ }
+
+ boolean hasSelectionQuery = selectionQuery != null &&
!selectionQuery.isBlank();
+ boolean hasMutationQuery = mutationQuery != null &&
!mutationQuery.isBlank();
+ boolean hasNaturalIdParameters = naturalIdParameters != null &&
!naturalIdParameters.isEmpty();
+ boolean hasStatelessOperation = statelessOperation != null &&
!statelessOperation.isBlank();
Review Comment:
💡 **Validation allows zero operations when none is configured but
`entityClassName` is set.** The validation counts `configuredOperations` and
requires exactly 1, which is correct. However, the endpoint has no way to
persist or merge an entity body — a common use case for a database component
producer. The only producer operations are query-based
(selectionQuery/mutationQuery) or bulk (statelessOperation insert/upsert). If a
user wants `from("direct:save").to("hibernate:MyEntity")` to persist the
exchange body, there's no option for that.
Consider whether this is intentional. If it is, document the limitation. If
not, consider adding a Session-based `persist`/`merge` operation.
##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateComponent.java:
##########
@@ -0,0 +1,238 @@
+/*
+ * 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.HashMap;
+import java.util.Map;
+
+import javax.sql.DataSource;
+
+import org.apache.camel.Endpoint;
+import org.apache.camel.RuntimeCamelException;
+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.hibernate.boot.MetadataSources;
+import org.hibernate.boot.registry.StandardServiceRegistry;
+import org.hibernate.boot.registry.StandardServiceRegistryBuilder;
+import org.hibernate.cfg.AvailableSettings;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+@Component("hibernate")
+public class HibernateComponent extends DefaultComponent {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(HibernateComponent.class);
+
+ @Metadata(description = "The Hibernate SessionFactory to use.")
+ private SessionFactory sessionFactory;
+
+ @Metadata(description = "The DataSource to use for bootstrapping the
SessionFactory (bean reference or registry name).")
+ private Object dataSource;
+
+ @Metadata(description = "Explicit array of entity classes.")
+ private Class[] entityClasses;
+
+ @Metadata(description = "Schema generation action: none, validate, update,
create.", defaultValue = "none")
+ private String schemaAction = "none";
+
+ @Metadata(description = "Arbitrary Hibernate configuration properties
passthrough map.")
Review Comment:
💡 **Raw type `Class[]`.** Should be `Class<?>[]` to avoid unchecked
warnings. Same for the getter/setter.
--
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]