This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-24159-eventhubs in repository https://gitbox.apache.org/repos/asf/camel.git
commit 4b18485fc796d32c7f7a54a2d0eb44bbc024949b Author: Claus Ibsen <[email protected]> AuthorDate: Fri Jul 17 19:52:22 2026 +0200 CAMEL-24159: camel-azure-eventhubs - fix medium-severity bugs from code review Co-Authored-By: Claude Opus 4.6 <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../component/azure/eventhubs/EventHubsCheckpointUpdaterTask.java | 4 ++-- .../camel/component/azure/eventhubs/EventHubsComponent.java | 8 +++++--- .../apache/camel/component/azure/eventhubs/EventHubsConsumer.java | 4 ++-- .../apache/camel/component/azure/eventhubs/EventHubsProducer.java | 8 +++++--- 4 files changed, 14 insertions(+), 10 deletions(-) diff --git a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsCheckpointUpdaterTask.java b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsCheckpointUpdaterTask.java index 49723ac7c4f7..e8a5355675d3 100644 --- a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsCheckpointUpdaterTask.java +++ b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsCheckpointUpdaterTask.java @@ -29,7 +29,7 @@ public class EventHubsCheckpointUpdaterTask implements Runnable { private static final Logger LOG = LoggerFactory.getLogger(EventHubsCheckpointUpdaterTask.class); - private EventContext eventContext; + private volatile EventContext eventContext; private final AtomicInteger processedEvents; private volatile long scheduledTime; @@ -44,7 +44,7 @@ public class EventHubsCheckpointUpdaterTask implements Runnable { LOG.debug("checkpointing offset after reaching timeout, with a batch of {}", processedEvents.get()); eventContext.updateCheckpointAsync() .subscribe(unused -> LOG.debug("Processed one event..."), - error -> LOG.debug("Error when updating Checkpoint: {}", error.getMessage()), + error -> LOG.warn("Error when updating Checkpoint: {}", error.getMessage(), error), () -> { LOG.debug("Checkpoint updated."); processedEvents.set(0); diff --git a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java index 04fd5b3de54e..489894ca9037 100644 --- a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java +++ b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java @@ -60,9 +60,11 @@ public class EventHubsComponent extends DefaultComponent { endpoint.getConfiguration().setCredentialType(CredentialType.CONNECTION_STRING); } } else { - boolean azure = endpoint.getConfiguration().getTokenCredential() instanceof DefaultAzureCredential; - endpoint.getConfiguration() - .setCredentialType(azure ? CredentialType.AZURE_IDENTITY : CredentialType.TOKEN_CREDENTIAL); + if (endpoint.getConfiguration().getCredentialType() == null) { + boolean azure = endpoint.getConfiguration().getTokenCredential() instanceof DefaultAzureCredential; + endpoint.getConfiguration() + .setCredentialType(azure ? CredentialType.AZURE_IDENTITY : CredentialType.TOKEN_CREDENTIAL); + } } validateConfigurations(configuration); diff --git a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java index f465326fbe86..b066caba6377 100644 --- a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java +++ b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java @@ -207,7 +207,7 @@ public class EventHubsConsumer extends DefaultConsumer implements ShutdownAware * * @param exchange the exchange */ - private void processCommit(final Exchange exchange, final EventContext eventContext) { + private synchronized void processCommit(final Exchange exchange, final EventContext eventContext) { if (lastTask == null || lastTask.isExpired()) { lastTask = new EventHubsCheckpointUpdaterTask(eventContext, processedEvents); // delegate the checkpoint update to a dedicated Thread @@ -230,7 +230,7 @@ public class EventHubsConsumer extends DefaultConsumer implements ShutdownAware } eventContext.updateCheckpointAsync() .subscribe(unused -> LOG.debug("Processed one event..."), - error -> LOG.debug("Error when updating Checkpoint: {}", error.getMessage()), + error -> LOG.warn("Error when updating Checkpoint: {}", error.getMessage(), error), () -> { LOG.debug("Checkpoint updated."); }); diff --git a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java index 861348cf52b5..2c5913a0bf75 100644 --- a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java +++ b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java @@ -27,6 +27,7 @@ import org.apache.camel.support.DefaultAsyncProducer; public class EventHubsProducer extends DefaultAsyncProducer { private EventHubProducerAsyncClient producerAsyncClient; + private boolean clientCloseable; private EventHubsProducerOperations producerOperations; public EventHubsProducer(final Endpoint endpoint) { @@ -40,8 +41,10 @@ public class EventHubsProducer extends DefaultAsyncProducer { EventHubsConfiguration configuration = getConfiguration(); producerAsyncClient = configuration.getProducerAsyncClient(); if (producerAsyncClient == null) { - // create the client producerAsyncClient = EventHubsClientFactory.createEventHubProducerAsyncClient(configuration); + clientCloseable = true; + } else { + clientCloseable = false; } // create our operations @@ -62,8 +65,7 @@ public class EventHubsProducer extends DefaultAsyncProducer { @Override protected void doStop() throws Exception { - if (producerAsyncClient != null) { - // shutdown async client + if (clientCloseable && producerAsyncClient != null) { producerAsyncClient.close(); producerAsyncClient = null; }
