This is an automated email from the ASF dual-hosted git repository.
gnodet pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 05127efd715a CAMEL-20199: Replace synchronized with ReentrantLock in
telemetry and metrics components (#25244)
05127efd715a is described below
commit 05127efd715aea64bd12c56002f5f4ccd69b1232
Author: Guillaume Nodet <[email protected]>
AuthorDate: Thu Jul 30 17:49:17 2026 +0200
CAMEL-20199: Replace synchronized with ReentrantLock in telemetry and
metrics components (#25244)
CAMEL-20199: Replace synchronized with ReentrantLock in telemetry and
metrics components
Convert synchronized blocks to ReentrantLock in camel-opentelemetry-metrics,
camel-opentelemetry2, and camel-telemetry-dev to eliminate virtual thread
pinning.
Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
.../opentelemetry/metrics/CounterProducer.java | 8 +++-
.../metrics/DistributionSummaryProducer.java | 8 +++-
.../camel/opentelemetry/metrics/TimerProducer.java | 8 +++-
.../OpenTelemetryRoutePolicyFactory.java | 21 ++++++----
.../camel/opentelemetry2/DevSpanExporter.java | 22 +++++++++--
.../camel/telemetrydev/InMemoryCollector.java | 45 ++++++++++++++--------
6 files changed, 79 insertions(+), 33 deletions(-)
diff --git
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/CounterProducer.java
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/CounterProducer.java
index a88a99c2f8d4..7bb72b1e60b0 100644
---
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/CounterProducer.java
+++
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/CounterProducer.java
@@ -18,6 +18,7 @@ package org.apache.camel.opentelemetry.metrics;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.LongUpDownCounter;
@@ -32,7 +33,7 @@ import static
org.apache.camel.opentelemetry.metrics.OpenTelemetryConstants.HEAD
public class CounterProducer extends
AbstractOpenTelemetryProducer<LongUpDownCounter> {
private final Map<String, LongUpDownCounter> counters = new
ConcurrentHashMap<>();
- private final Object lock = new Object();
+ private final ReentrantLock lock = new ReentrantLock();
public CounterProducer(OpenTelemetryEndpoint endpoint) {
super(endpoint);
@@ -42,7 +43,8 @@ public class CounterProducer extends
AbstractOpenTelemetryProducer<LongUpDownCou
protected LongUpDownCounter getInstrument(String name, String description)
{
LongUpDownCounter counter = counters.get(name);
if (counter == null) {
- synchronized (lock) {
+ lock.lock();
+ try {
counter = counters.get(name);
if (counter == null) {
Meter meter = getEndpoint().getMeter();
@@ -53,6 +55,8 @@ public class CounterProducer extends
AbstractOpenTelemetryProducer<LongUpDownCou
counter = builder.build();
counters.put(name, counter);
}
+ } finally {
+ lock.unlock();
}
}
return counter;
diff --git
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/DistributionSummaryProducer.java
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/DistributionSummaryProducer.java
index c3b5a2bb4bc3..6544c168ed70 100644
---
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/DistributionSummaryProducer.java
+++
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/DistributionSummaryProducer.java
@@ -18,6 +18,7 @@ package org.apache.camel.opentelemetry.metrics;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.LongHistogram;
@@ -30,7 +31,7 @@ import static
org.apache.camel.opentelemetry.metrics.OpenTelemetryConstants.HEAD
public class DistributionSummaryProducer extends
AbstractOpenTelemetryProducer<LongHistogram> {
private final Map<String, LongHistogram> distributionSummaries = new
ConcurrentHashMap<>();
- private final Object lock = new Object();
+ private final ReentrantLock lock = new ReentrantLock();
public DistributionSummaryProducer(OpenTelemetryEndpoint endpoint) {
super(endpoint);
@@ -40,7 +41,8 @@ public class DistributionSummaryProducer extends
AbstractOpenTelemetryProducer<L
protected LongHistogram getInstrument(String name, String description) {
LongHistogram summary = distributionSummaries.get(name);
if (summary == null) {
- synchronized (lock) {
+ lock.lock();
+ try {
summary = distributionSummaries.get(name);
if (summary == null) {
Meter meter = getEndpoint().getMeter();
@@ -51,6 +53,8 @@ public class DistributionSummaryProducer extends
AbstractOpenTelemetryProducer<L
summary = builder.build();
distributionSummaries.put(name, summary);
}
+ } finally {
+ lock.unlock();
}
}
return summary;
diff --git
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/TimerProducer.java
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/TimerProducer.java
index b9cf907700dd..396a9f2302ba 100644
---
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/TimerProducer.java
+++
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/TimerProducer.java
@@ -18,6 +18,7 @@ package org.apache.camel.opentelemetry.metrics;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.LongHistogram;
@@ -31,7 +32,7 @@ import static
org.apache.camel.opentelemetry.metrics.OpenTelemetryConstants.HEAD
public class TimerProducer extends
AbstractOpenTelemetryProducer<LongHistogram> {
private final Map<String, LongHistogram> timers = new
ConcurrentHashMap<>();
- private final Object lock = new Object();
+ private final ReentrantLock lock = new ReentrantLock();
public TimerProducer(OpenTelemetryEndpoint endpoint) {
super(endpoint);
@@ -41,7 +42,8 @@ public class TimerProducer extends
AbstractOpenTelemetryProducer<LongHistogram>
protected LongHistogram getInstrument(String name, String description) {
LongHistogram timer = timers.get(name);
if (timer == null) {
- synchronized (lock) {
+ lock.lock();
+ try {
timer = timers.get(name);
if (timer == null) {
Meter meter = getEndpoint().getMeter();
@@ -53,6 +55,8 @@ public class TimerProducer extends
AbstractOpenTelemetryProducer<LongHistogram>
timer = builder.build();
timers.put(name, timer);
}
+ } finally {
+ lock.unlock();
}
}
return timer;
diff --git
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/routepolicy/OpenTelemetryRoutePolicyFactory.java
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/routepolicy/OpenTelemetryRoutePolicyFactory.java
index e3fbef79681e..67ffe826b87f 100644
---
a/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/routepolicy/OpenTelemetryRoutePolicyFactory.java
+++
b/components/camel-opentelemetry-metrics/src/main/java/org/apache/camel/opentelemetry/metrics/routepolicy/OpenTelemetryRoutePolicyFactory.java
@@ -17,6 +17,7 @@
package org.apache.camel.opentelemetry.metrics.routepolicy;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.locks.ReentrantLock;
import io.opentelemetry.api.metrics.Meter;
import org.apache.camel.CamelContext;
@@ -34,6 +35,7 @@ import org.apache.camel.support.service.ServiceSupport;
public class OpenTelemetryRoutePolicyFactory extends ServiceSupport
implements RoutePolicyFactory, CamelContextAware, NonManagedService,
StaticService {
+ private final ReentrantLock lock = new ReentrantLock();
private CamelContext camelContext;
private Meter meter;
private RouteMetric contextMetric;
@@ -92,14 +94,19 @@ public class OpenTelemetryRoutePolicyFactory extends
ServiceSupport
this.longTaskTimeUnit = longTaskTimeUnit;
}
- public synchronized RouteMetric
createOrGetContextMetric(OpenTelemetryRoutePolicy policy) {
- if (contextMetric == null) {
- contextMetric = new OpenTelemetryContextMetricsStatistics(
- meter, camelContext, policy.getNamingStrategy(),
policy.getConfiguration(),
- policy.isRegisterKamelets(), policy.isRegisterTemplates(),
- policy.getTimeUnit(), policy.getLongTaskTimeUnit());
+ public RouteMetric createOrGetContextMetric(OpenTelemetryRoutePolicy
policy) {
+ lock.lock();
+ try {
+ if (contextMetric == null) {
+ contextMetric = new OpenTelemetryContextMetricsStatistics(
+ meter, camelContext, policy.getNamingStrategy(),
policy.getConfiguration(),
+ policy.isRegisterKamelets(),
policy.isRegisterTemplates(),
+ policy.getTimeUnit(), policy.getLongTaskTimeUnit());
+ }
+ return contextMetric;
+ } finally {
+ lock.unlock();
}
- return contextMetric;
}
@Override
diff --git
a/components/camel-opentelemetry2/src/main/java/org/apache/camel/opentelemetry2/DevSpanExporter.java
b/components/camel-opentelemetry2/src/main/java/org/apache/camel/opentelemetry2/DevSpanExporter.java
index 5767e9f76f3b..416e9451be37 100644
---
a/components/camel-opentelemetry2/src/main/java/org/apache/camel/opentelemetry2/DevSpanExporter.java
+++
b/components/camel-opentelemetry2/src/main/java/org/apache/camel/opentelemetry2/DevSpanExporter.java
@@ -21,6 +21,7 @@ import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.locks.ReentrantLock;
import io.opentelemetry.sdk.common.CompletableResultCode;
import io.opentelemetry.sdk.trace.data.SpanData;
@@ -39,6 +40,7 @@ final class DevSpanExporter implements SpanExporter {
static final int DEFAULT_CAPACITY = 2000;
+ private final ReentrantLock lock = new ReentrantLock();
private final int capacity;
private volatile boolean stopped;
@@ -59,7 +61,8 @@ final class DevSpanExporter implements SpanExporter {
if (stopped) {
return CompletableResultCode.ofSuccess();
}
- synchronized (traces) {
+ lock.lock();
+ try {
for (SpanData span : spanDataList) {
String traceId = span.getTraceId();
traces.computeIfAbsent(traceId, k -> new
ArrayList<>()).add(span);
@@ -74,6 +77,8 @@ final class DevSpanExporter implements SpanExporter {
it.remove();
}
}
+ } finally {
+ lock.unlock();
}
return CompletableResultCode.ofSuccess();
}
@@ -90,18 +95,24 @@ final class DevSpanExporter implements SpanExporter {
}
List<SpanData> getFinishedSpans() {
- synchronized (traces) {
+ lock.lock();
+ try {
List<SpanData> result = new ArrayList<>(totalSpanCount);
for (List<SpanData> traceSpans : traces.values()) {
result.addAll(traceSpans);
}
return result;
+ } finally {
+ lock.unlock();
}
}
int getSpanCount() {
- synchronized (traces) {
+ lock.lock();
+ try {
return totalSpanCount;
+ } finally {
+ lock.unlock();
}
}
@@ -110,9 +121,12 @@ final class DevSpanExporter implements SpanExporter {
}
void reset() {
- synchronized (traces) {
+ lock.lock();
+ try {
traces.clear();
totalSpanCount = 0;
+ } finally {
+ lock.unlock();
}
}
}
diff --git
a/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/InMemoryCollector.java
b/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/InMemoryCollector.java
index 8baad8bf54d9..f51040cf67a8 100644
---
a/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/InMemoryCollector.java
+++
b/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/InMemoryCollector.java
@@ -19,36 +19,49 @@ package org.apache.camel.telemetrydev;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.Map;
+import java.util.concurrent.locks.ReentrantLock;
/*
* Basic inmemory implementation for an home made spans collector.
*/
public class InMemoryCollector {
+ private final ReentrantLock lock = new ReentrantLock();
+
// traceid --> spanid --> Span
private Map<String, Map<String, DevSpanAdapter>> traceDB = new HashMap<>();
- public synchronized void push(String traceId, DevSpanAdapter span) {
- Map<String, DevSpanAdapter> spans = traceDB.get(traceId);
- if (spans == null) {
- spans = new HashMap<>();
- traceDB.put(traceId, spans);
+ public void push(String traceId, DevSpanAdapter span) {
+ lock.lock();
+ try {
+ Map<String, DevSpanAdapter> spans = traceDB.get(traceId);
+ if (spans == null) {
+ spans = new HashMap<>();
+ traceDB.put(traceId, spans);
+ }
+ spans.put(span.getTag("spanid"), span);
+ } finally {
+ lock.unlock();
}
- spans.put(span.getTag("spanid"), span);
}
- public synchronized DevTrace get(String traceId) {
- Map<String, DevSpanAdapter> spans = traceDB.get(traceId);
- if (spans == null) {
- return null;
- }
- for (DevSpanAdapter span : spans.values()) {
- if (!"true".equals(span.getTag("isDone"))) {
- // Still an active trace, not all spans are closed
+ public DevTrace get(String traceId) {
+ lock.lock();
+ try {
+ Map<String, DevSpanAdapter> spans = traceDB.get(traceId);
+ if (spans == null) {
return null;
}
+ for (DevSpanAdapter span : spans.values()) {
+ if (!"true".equals(span.getTag("isDone"))) {
+ // Still an active trace, not all spans are closed
+ return null;
+ }
+ }
+ traceDB.remove(traceId);
+ return new DevTrace(traceId, new ArrayList<>(spans.values()));
+ } finally {
+ lock.unlock();
}
- traceDB.remove(traceId);
- return new DevTrace(traceId, new ArrayList<>(spans.values()));
}
}