Асинхронная публикация, игнорирующая ошибки и таймаутыJAVA

Программисты JAVA общаются здесь
Anonymous
Асинхронная публикация, игнорирующая ошибки и таймауты

Сообщение Anonymous »

У меня есть приложение Spring Boot, которое отправляет через RabbitMq статистику ответа на каждый запрос контроллера. Он работает, когда сервер RabbitMq доступен и принимает сообщение.
Однако эта публикация должна быть необязательной, и на клиент не должно влиять, если по какой-либо причине RabbitMq недоступен или не может подтвердите сообщение. Я хотел бы иметь тихий сбой, не влияющий на время ответа и задержку REST API (только регистрация ошибок). К сожалению, я наблюдаю блокирующий вызов в AsyncRabbitTemplate::convertSendAndReceive в случае таймаута соединения с сервером RabbitMq.
Мы используем Java 21, поэтому в качестве обходного пути я решил использовать исполнитель с виртуальными потоками, который разблокирует публикацию, например:
@Component
public class RabbitMQProducer {
private static final Logger LOG = LoggerFactory.getLogger(RabbitMQProducer.class);
private static final RateLimitedLog RATE_LIMITED_LOG = RateLimitedLog.withRateLimit(LOG).maxRate(10).every(Duration.ofSeconds(15)).build();

private final AsyncRabbitTemplate rabbitTemplate;

private final Executor executor = Executors.newVirtualThreadPerTaskExecutor();

public RabbitMQProducer(AsyncRabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}

public void sendMessage(String exchangeName, String routingKey, String message) {
executor.execute(() -> send(exchangeName, routingKey, message));
}

private void send(String exchangeName, String routingKey, String message) {
final RabbitConverterFuture future =
rabbitTemplate.convertSendAndReceive(exchangeName, routingKey, message);
future.whenComplete((result, ex) -> {
if (ex != null) {
RATE_LIMITED_LOG.error(message, ex);
}
});
}
}

Кажется, он работает нормально, если RabbitMq не отвечает (мы развертываем наше приложение в Kubernetes и проверяем, что я только что отключил сетевое подключение путем добавления NetworkPolicies, запрещающего подключение).
Так вот вопрос: зачем нужен исполнитель с Virtual Thread? Если я его удалю:
public void sendMessage(String exchangeName, String routingKey, String message) {
final RabbitConverterFuture future =
rabbitTemplate.convertSendAndReceive(exchangeName, routingKey, message);
future.whenComplete((result, ex) -> {
if (ex != null) {
RATE_LIMITED_LOG.error(message, ex);
}
});
}
}

метод фактически блокируется, и это влияет на конечных пользователей.
Вот конфигурационные компоненты RabbitMq:
< pre class="lang-java Prettyprint-override">@Configuration
public class RabbitMQConfig {
@Bean
public ConnectionFactory connectionFactory(RabbitMqConfiguration config) {
final CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory();

cachingConnectionFactory.setHost(config.getAmqpHost());
cachingConnectionFactory.setPort(config.getAmqpPort());
cachingConnectionFactory.setUsername(config.getAmqpUser());
cachingConnectionFactory.setPassword(config.getAmqpPassword());

cachingConnectionFactory.getRabbitConnectionFactory().setAutomaticRecoveryEnabled(true);

return cachingConnectionFactory;
}

@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
final RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);

rabbitTemplate.setRetryTemplate(RetryTemplate.builder().maxAttempts(1).build());

return rabbitTemplate;
}

@Bean
public AsyncRabbitTemplate asyncRabbitTemplate(RabbitTemplate rabbitTemplate) {
return new AsyncRabbitTemplate(rabbitTemplate);
}

@Bean
public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(
SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory connectionFactory) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
configurer.configure(factory, connectionFactory);
factory.setConsumerTagStrategy(q -> CONSUMER_TAG);
return factory;
}

@Bean(name = "searchStreamExchange")
public Exchange searchStreamExchange(RabbitMqConfiguration config) {
return new TopicExchange(config.getSearchStreamExchange(), false, false);
}
}



Подробнее здесь: https://stackoverflow.com/questions/789 ... d-timeouts

Вернуться в «JAVA»