Change "synchronized" to reentrant lock for virtual-threads
Fix checkstyles before merge Code cleanup Double-Checked Locking Optimization was used to avoid unnecessary locking overhead.
This commit is contained in:
committed by
Oleg Zhurakousky
parent
7a5e5d0541
commit
cbfd3aa995
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2019-2019 the original author or authors.
|
||||
* Copyright 2019-2024 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -24,6 +24,7 @@ import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.ToDoubleFunction;
|
||||
|
||||
import io.micrometer.core.instrument.FunctionCounter;
|
||||
@@ -52,6 +53,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean;
|
||||
* We will keep this class, as long as we support Boot 2.2.x.
|
||||
*
|
||||
* @author Soby Chacko
|
||||
* @author Omer Celik
|
||||
* @since 3.0.0
|
||||
*/
|
||||
public class KafkaStreamsBinderMetrics {
|
||||
@@ -84,6 +86,8 @@ public class KafkaStreamsBinderMetrics {
|
||||
|
||||
private volatile Set<MetricName> currentMeters = new HashSet<>();
|
||||
|
||||
private static final ReentrantLock metricsLock = new ReentrantLock();
|
||||
|
||||
public KafkaStreamsBinderMetrics(MeterRegistry meterRegistry) {
|
||||
this.meterRegistry = meterRegistry;
|
||||
}
|
||||
@@ -108,9 +112,13 @@ public class KafkaStreamsBinderMetrics {
|
||||
}
|
||||
|
||||
public void addMetrics(Set<StreamsBuilderFactoryBean> streamsBuilderFactoryBeans) {
|
||||
synchronized (KafkaStreamsBinderMetrics.this) {
|
||||
try {
|
||||
metricsLock.lock();
|
||||
this.bindTo(streamsBuilderFactoryBeans);
|
||||
}
|
||||
finally {
|
||||
metricsLock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
void prepareToBindMetrics(MeterRegistry registry, Map<MetricName, ? extends Metric> metrics) {
|
||||
|
||||
@@ -30,6 +30,7 @@ import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.ToDoubleFunction;
|
||||
|
||||
import io.micrometer.core.instrument.Gauge;
|
||||
@@ -68,6 +69,7 @@ import org.springframework.util.ObjectUtils;
|
||||
* @author Tomek Szmytka
|
||||
* @author Nico Heller
|
||||
* @author Kurt Hong
|
||||
* @author Omer Celik
|
||||
*/
|
||||
public class KafkaBinderMetrics
|
||||
implements MeterBinder, ApplicationListener<BindingCreatedEvent>, AutoCloseable {
|
||||
@@ -97,6 +99,8 @@ public class KafkaBinderMetrics
|
||||
|
||||
ScheduledExecutorService scheduler;
|
||||
|
||||
private final ReentrantLock consumerFactoryLock = new ReentrantLock();
|
||||
|
||||
public KafkaBinderMetrics(KafkaMessageChannelBinder binder,
|
||||
KafkaBinderConfigurationProperties binderConfigurationProperties,
|
||||
ConsumerFactory<?, ?> defaultConsumerFactory,
|
||||
@@ -231,25 +235,36 @@ public class KafkaBinderMetrics
|
||||
return lag;
|
||||
}
|
||||
|
||||
private synchronized ConsumerFactory<?, ?> createConsumerFactory() {
|
||||
/**
|
||||
* Double-Checked Locking Optimization was used to avoid unnecessary locking overhead.
|
||||
*/
|
||||
private ConsumerFactory<?, ?> createConsumerFactory() {
|
||||
if (this.defaultConsumerFactory == null) {
|
||||
Map<String, Object> props = new HashMap<>();
|
||||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
|
||||
ByteArrayDeserializer.class);
|
||||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
|
||||
ByteArrayDeserializer.class);
|
||||
Map<String, Object> mergedConfig = this.binderConfigurationProperties
|
||||
.mergedConsumerConfiguration();
|
||||
if (!ObjectUtils.isEmpty(mergedConfig)) {
|
||||
props.putAll(mergedConfig);
|
||||
}
|
||||
if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) {
|
||||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
|
||||
this.binderConfigurationProperties
|
||||
try {
|
||||
this.consumerFactoryLock.lock();
|
||||
if (this.defaultConsumerFactory == null) {
|
||||
Map<String, Object> props = new HashMap<>();
|
||||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
|
||||
ByteArrayDeserializer.class);
|
||||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
|
||||
ByteArrayDeserializer.class);
|
||||
Map<String, Object> mergedConfig = this.binderConfigurationProperties
|
||||
.mergedConsumerConfiguration();
|
||||
if (!ObjectUtils.isEmpty(mergedConfig)) {
|
||||
props.putAll(mergedConfig);
|
||||
}
|
||||
if (!props.containsKey(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)) {
|
||||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
|
||||
this.binderConfigurationProperties
|
||||
.getKafkaConnectionString());
|
||||
}
|
||||
this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>(
|
||||
props);
|
||||
}
|
||||
}
|
||||
finally {
|
||||
this.consumerFactoryLock.unlock();
|
||||
}
|
||||
this.defaultConsumerFactory = new DefaultKafkaConsumerFactory<>(
|
||||
props);
|
||||
}
|
||||
return this.defaultConsumerFactory;
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2023 the original author or authors.
|
||||
* Copyright 2016-2024 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -22,6 +22,7 @@ import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.IntStream;
|
||||
import java.util.stream.Stream;
|
||||
@@ -79,6 +80,7 @@ import org.springframework.util.StringUtils;
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Michael Michailidis
|
||||
* @author Byungjun You
|
||||
* @author Omer Celik
|
||||
*/
|
||||
// @checkstyle:off
|
||||
public class RabbitExchangeQueueProvisioner
|
||||
@@ -105,6 +107,8 @@ public class RabbitExchangeQueueProvisioner
|
||||
|
||||
private final AtomicInteger producerExchangeBeanNameQualifier = new AtomicInteger();
|
||||
|
||||
private final ReentrantLock autoDeclareContextLock = new ReentrantLock();
|
||||
|
||||
public RabbitExchangeQueueProvisioner(ConnectionFactory connectionFactory) {
|
||||
this(connectionFactory, Collections.emptyList());
|
||||
}
|
||||
@@ -320,11 +324,15 @@ public class RabbitExchangeQueueProvisioner
|
||||
(q, i) -> IntStream.range(0, i)
|
||||
.mapToObj(j -> rk + "-" + j)
|
||||
.collect(Collectors.toList()));
|
||||
synchronized (this.autoDeclareContext) {
|
||||
try {
|
||||
this.autoDeclareContextLock.lock();
|
||||
if (!this.autoDeclareContext.containsBean(name + ".superStream")) {
|
||||
this.autoDeclareContext.getBeanFactory().registerSingleton(name + ".superStream", ss);
|
||||
}
|
||||
}
|
||||
finally {
|
||||
this.autoDeclareContextLock.unlock();
|
||||
}
|
||||
try {
|
||||
ss.getDeclarables().forEach(dec -> {
|
||||
if (dec instanceof Exchange exch) {
|
||||
@@ -716,7 +724,8 @@ public class RabbitExchangeQueueProvisioner
|
||||
}
|
||||
|
||||
private void addToAutoDeclareContext(String name, Declarable bean) {
|
||||
synchronized (this.autoDeclareContext) {
|
||||
try {
|
||||
this.autoDeclareContextLock.lock();
|
||||
if (!this.autoDeclareContext.containsBean(name)) {
|
||||
this.autoDeclareContext.getBeanFactory().registerSingleton(name, new Declarables(bean));
|
||||
}
|
||||
@@ -724,6 +733,9 @@ public class RabbitExchangeQueueProvisioner
|
||||
this.autoDeclareContext.getBean(name, Declarables.class).getDeclarables().add(bean);
|
||||
}
|
||||
}
|
||||
finally {
|
||||
this.autoDeclareContextLock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
private void declareBinding(String rootName, org.springframework.amqp.core.Binding bindingArg) {
|
||||
@@ -756,31 +768,36 @@ public class RabbitExchangeQueueProvisioner
|
||||
public void cleanAutoDeclareContext(ConsumerDestination destination,
|
||||
ExtendedConsumerProperties<RabbitConsumerProperties> consumerProperties) {
|
||||
|
||||
synchronized (this.autoDeclareContext) {
|
||||
try {
|
||||
this.autoDeclareContextLock.lock();
|
||||
Stream.of(StringUtils.tokenizeToStringArray(destination.getName(), ",", true,
|
||||
true)).forEach(name -> {
|
||||
String group = null;
|
||||
String bindingName = null;
|
||||
if (destination instanceof RabbitConsumerDestination rabbitConsumerDestination) {
|
||||
group = rabbitConsumerDestination.getGroup();
|
||||
bindingName = rabbitConsumerDestination.getBindingName();
|
||||
}
|
||||
RabbitConsumerProperties properties = consumerProperties.getExtension();
|
||||
String toRemove = properties.isQueueNameGroupOnly() ? bindingName + "." + group : name.trim();
|
||||
boolean partitioned = consumerProperties.isPartitioned();
|
||||
if (partitioned) {
|
||||
toRemove = removePartitionPart(toRemove);
|
||||
}
|
||||
removeSingleton(toRemove + ".exchange");
|
||||
removeQueueAndBindingBeans(properties, name.trim(), "", group, partitioned);
|
||||
});
|
||||
true)).forEach(name -> {
|
||||
String group = null;
|
||||
String bindingName = null;
|
||||
if (destination instanceof RabbitConsumerDestination rabbitConsumerDestination) {
|
||||
group = rabbitConsumerDestination.getGroup();
|
||||
bindingName = rabbitConsumerDestination.getBindingName();
|
||||
}
|
||||
RabbitConsumerProperties properties = consumerProperties.getExtension();
|
||||
String toRemove = properties.isQueueNameGroupOnly() ? bindingName + "." + group : name.trim();
|
||||
boolean partitioned = consumerProperties.isPartitioned();
|
||||
if (partitioned) {
|
||||
toRemove = removePartitionPart(toRemove);
|
||||
}
|
||||
removeSingleton(toRemove + ".exchange");
|
||||
removeQueueAndBindingBeans(properties, name.trim(), "", group, partitioned);
|
||||
});
|
||||
}
|
||||
finally {
|
||||
this.autoDeclareContextLock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
public void cleanAutoDeclareContext(ProducerDestination dest,
|
||||
ExtendedProducerProperties<RabbitProducerProperties> properties) {
|
||||
|
||||
synchronized (this.autoDeclareContext) {
|
||||
try {
|
||||
this.autoDeclareContextLock.lock();
|
||||
if (dest instanceof RabbitProducerDestination rabbitProducerDestination) {
|
||||
String qual = rabbitProducerDestination.getBeanNameQualifier();
|
||||
removeSingleton(dest.getName() + "." + qual + ".exchange");
|
||||
@@ -790,13 +807,13 @@ public class RabbitExchangeQueueProvisioner
|
||||
if (properties.isPartitioned()) {
|
||||
for (int i = 0; i < properties.getPartitionCount(); i++) {
|
||||
removeQueueAndBindingBeans(properties.getExtension(),
|
||||
properties.getExtension().isQueueNameGroupOnly() ? "" : dest.getName(),
|
||||
group + "-" + i, group, true);
|
||||
properties.getExtension().isQueueNameGroupOnly() ? "" : dest.getName(),
|
||||
group + "-" + i, group, true);
|
||||
}
|
||||
}
|
||||
else {
|
||||
removeQueueAndBindingBeans(properties.getExtension(), dest.getName() + "." + group, "",
|
||||
group, false);
|
||||
group, false);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -811,6 +828,9 @@ public class RabbitExchangeQueueProvisioner
|
||||
}
|
||||
}
|
||||
}
|
||||
finally {
|
||||
this.autoDeclareContextLock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
private void removeQueueAndBindingBeans(RabbitCommonProperties properties, String name, String suffix,
|
||||
|
||||
Reference in New Issue
Block a user