Java for Beginner
869 subscribers
1.01K photos
275 videos
14 files
1.69K links
Канал от новичков для новичков!
Изучайте Java вместе с нами!
Здесь мы обмениваемся опытом и постоянно изучаем что-то новое!

Наш YouTube канал - https://www.youtube.com/@Java_Beginner-Dev

Наш канал на RUTube - https://rutube.ru/channel/37896292/
Download Telegram
Spring Boot 3.x + RabbitMQ: Продвинутая конфигурация

Введение: Современный стек для production

Spring Boot 3.x в сочетании с RabbitMQ представляет собой enterprise-уровневую платформу для построения event-driven систем. В 2025 году фокус сместился с базовой интеграции на production-готовые конфигурации, обеспечивающие безопасность, наблюдаемость и отказоустойчивость.


Spring Boot 3.3+ и RabbitMQ Starter: ключевые изменения

Spring Boot 3.3 существенно улучшил автоконфигурацию RabbitMQ:
Автоматическое определение кворумных очередей (Quorum Queues) как стандарта
Интеграция с Micrometer для детального мониторинга
Поддержка реактивных потоков через Project Reactor
Готовность к GraalVM native images из коробки

Минимальная конфигурация:
@Configuration
public class RabbitMQConfig {

@Bean
public CachingConnectionFactory connectionFactory(
@Value("${spring.rabbitmq.host}") String host) {
CachingConnectionFactory factory = new CachingConnectionFactory(host);

// Производственные настройки по умолчанию
factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
factory.setPublisherReturns(true);
factory.setChannelCacheSize(25);

return factory;
}
}



Конфигурация для production: критичные параметры

Для production обязательны:
TLS/SSL для шифрования трафика
mTLS (Mutual TLS) для аутентификации клиентов
OAuth2 или JWT для авторизации

Конфигурация mTLS в application.yaml:
spring:
rabbitmq:
ssl:
enabled: true
algorithm: TLSv1.3
key-store: file:/etc/ssl/client-keystore.p12
key-store-password: ${KEYSTORE_PASSWORD}
key-store-type: PKCS12
trust-store: file:/etc/ssl/truststore.jks
trust-store-password: ${TRUSTSTORE_PASSWORD}
trust-store-type: JKS
verify-hostname: true


Гарантии доставки
spring:
rabbitmq:
# Подтверждения от брокера
publisher-confirm-type: correlated
publisher-returns: true

# Подтверждения от потребителя
listener:
simple:
acknowledge-mode: manual # Для полного контроля
prefetch: 10 # Балансировка нагрузки
retry:
enabled: true
max-attempts: 3
back-off:
initial-interval: 1s
multiplier: 2


Spring Boot 3.3 предоставляет расширенный мониторинг через Micrometer:

@Configuration
public class MonitoringConfig {

@Bean
public MeterRegistryCustomizer<MeterRegistry> rabbitMQMetrics() {
return registry -> {
// Кастомные метрики для RabbitMQ
registry.config().commonTags("application", "order-service");
};
}

@Bean
public RabbitTemplateMonitoring rabbitTemplateMonitoring(
RabbitTemplate rabbitTemplate,
MeterRegistry meterRegistry) {

// Мониторинг времени доставки
Timer timer = Timer.builder("rabbitmq.publish.duration")
.description("Time to publish message to RabbitMQ")
.register(meterRegistry);

rabbitTemplate.setObservationEnabled(true);

return new RabbitTemplateMonitoring(rabbitTemplate, timer);
}
}



Reactive RabbitMQ с Project Reactor

Реактивная модель идеальна для:
Высоконагруженных систем с тысячами соединений
Сценариев, требующих backpressure управления
Интеграции с другими реактивными компонентами (WebFlux, R2DBC)

#Java #middle #RabbitMQ
👍2🔥1
Конфигурация реактивного RabbitMQ:
@Configuration
@EnableReactiveRabbit
public class ReactiveRabbitConfig {

@Bean
public ReactiveRabbitTemplate reactiveRabbitTemplate(
ReactiveRabbitConnectionFactory factory) {

ReactiveRabbitTemplate template = new ReactiveRabbitTemplate(factory);

// Настройки backpressure
template.setObservationEnabled(true);
template.setUsePublisherConfirm(true);

return template;
}

@Bean
public ReactiveRabbitListenerContainerFactory containerFactory(
ReactiveRabbitConnectionFactory factory) {

ReactiveRabbitListenerContainerFactory containerFactory =
new ReactiveRabbitListenerContainerFactory();
containerFactory.setConnectionFactory(factory);

// Обработка ошибок в реактивном стиле
containerFactory.setErrorHandler(errorHandler());

return containerFactory;
}

private ReactiveErrorHandler errorHandler() {
return new ReactiveErrorHandler() {
@Override
public Mono<Void> handleError(Throwable throwable,
Acknowledgement acknowledgement) {
log.error("Reactive consumption error", throwable);
return acknowledgement.nack(false); // Не requeue при фатальных ошибках
}
};
}
}


Реактивный consumer с обработкой backpressure:
@Component
public class OrderEventConsumer {

@ReactiveRabbitListener(queues = "orders.queue")
public Mono<Void> processOrder(OrderEvent event,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {

return Mono.fromCallable(() -> {
// Обработка события
return processOrderInternal(event);
})
.doOnSuccess(result -> {
// Ручное подтверждение
channel.basicAck(tag, false);
})
.doOnError(error -> {
// Логика повторной обработки
channel.basicNack(tag, false, shouldRequeue(error));
})
.then();
}
}



GraalVM native поддержка


Spring Boot 3.3 значительно улучшил поддержку GraalVM native images:

Необходимые зависимости:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-aot</artifactId>
</dependency>

<dependency>
<groupId>org.graalvm.buildtools</groupId>
<artifactId>native-maven-plugin</artifactId>
</dependency>


Конфигурация для RabbitMQ в native-образе:
@NativeHint(
types = {
@TypeHint(types = {
com.rabbitmq.client.impl.AMQConnection.class,
com.rabbitmq.client.impl.recovery.AutorecoveringConnection.class,
// Критичные классы для RabbitMQ клиента
})
},
initialization = {
@InitializationHint(
types = com.rabbitmq.client.impl.AMQConnection.class,
initTime = InitializationTime.BUILD
)
}
)
@Configuration
public class RabbitMQNativeConfig {

// Явное объявление бинов, необходимых в runtime
@Bean
@RuntimeHint
public ConnectionFactory connectionFactory() {
// Фабрика должна создаваться с явными типами
return new CachingConnectionFactory();
}
}


Проблемы и решения для native-сборок:
Reflection: RabbitMQ клиент использует reflection для создания соединений
Динамическая загрузка классов: Требует явного конфигурирования в native-image
Сетевые операции: Нативный образ требует явного разрешения сетевых операций


#Java #middle #RabbitMQ
👍3🔥1
Контейнеризация: health checks и liveness probes

Spring Boot Actuator предоставляет расширенные health indicators для RabbitMQ:

management:
endpoints:
web:
exposure:
include: health,metrics,prometheus

endpoint:
health:
show-details: always
probes:
enabled: true

health:
rabbit:
enabled: true

livenessstate:
enabled: true
readinessstate:
enabled: true


Кастомный health indicator для бизнес-логики:
@Component
public class RabbitMQBusinessHealthIndicator
implements HealthIndicator {

private final RabbitAdmin rabbitAdmin;
private final MeterRegistry meterRegistry;

public RabbitMQBusinessHealthIndicator(RabbitAdmin rabbitAdmin,
MeterRegistry meterRegistry) {
this.rabbitAdmin = rabbitAdmin;
this.meterRegistry = meterRegistry;
}

@Override
public Health health() {
try {
// Проверка критичных очередей
Properties queueProperties = rabbitAdmin.getQueueProperties("orders.queue");

if (queueProperties == null) {
return Health.down()
.withDetail("error", "Critical queue not found")
.build();
}

int messageCount = (int) queueProperties.get("QUEUE_MESSAGE_COUNT");

// Проверка на переполнение
if (messageCount > 10000) {
return Health.down()
.withDetail("warning", "Queue backlog too high")
.withDetail("messageCount", messageCount)
.build();
}

return Health.up()
.withDetail("queueStatus", "healthy")
.withDetail("messageCount", messageCount)
.build();

} catch (Exception e) {
return Health.down(e).build();
}
}
}


Kubernetes liveness и readiness probes

# Kubernetes deployment конфигурация
apiVersion: apps/v1
kind: Deployment
spec:
template:
spec:
containers:
- name: order-service
livenessProbe:
httpGet:
path: /actuator/health/liveness
port: 8080
initialDelaySeconds: 90 # Учитываем старт Spring Boot 3.x
periodSeconds: 15
failureThreshold: 3

readinessProbe:
httpGet:
path: /actuator/health/readiness
port: 8080
initialDelaySeconds: 30
periodSeconds: 5
successThreshold: 2
failureThreshold: 5

startupProbe:
httpGet:
path: /actuator/health/startup
port: 8080
failureThreshold: 30 # Долгий старт допустим
periodSeconds: 10



Security: продвинутые сценарии

OAuth2 с Resource Server
@Configuration
@EnableWebSecurity
public class RabbitMQSecurityConfig {

@Bean
public SecurityFilterChain securityFilterChain(HttpSecurity http) throws Exception {
http
.authorizeHttpRequests(authz -> authz
.requestMatchers("/actuator/**").permitAll()
.anyRequest().authenticated()
)
.oauth2ResourceServer(oauth2 -> oauth2
.jwt(jwt -> jwt
.jwtAuthenticationConverter(jwtAuthenticationConverter())
)
);

return http.build();
}

@Bean
public JwtAuthenticationConverter jwtAuthenticationConverter() {
JwtGrantedAuthoritiesConverter converter = new JwtGrantedAuthoritiesConverter();
converter.setAuthorityPrefix("RABBITMQ_");
converter.setAuthoritiesClaimName("rabbitmq_roles");

JwtAuthenticationConverter jwtConverter = new JwtAuthenticationConverter();
jwtConverter.setJwtGrantedAuthoritiesConverter(converter);

return jwtConverter;
}
}


#Java #middle #RabbitMQ
👍2🔥1
Интеграция OAuth2 с RabbitTemplate:
@Component
public class OAuth2RabbitTemplateDecorator {

private final ReactiveClientRegistrationRepository clientRegistrationRepository;

public OAuth2RabbitTemplateDecorator(
ReactiveClientRegistrationRepository clientRegistrationRepository) {
this.clientRegistrationRepository = clientRegistrationRepository;
}

public Mono<RabbitTemplate> decorateWithOAuth2(RabbitTemplate template) {
return clientRegistrationRepository
.findByRegistrationId("rabbitmq")
.flatMap(clientRegistration -> {
// Получение OAuth2 токена для RabbitMQ
return retrieveAccessToken(clientRegistration);
})
.map(accessToken -> {
// Декорация сообщений с токеном
template.addBeforePublishPostProcessors(message -> {
message.getMessageProperties()
.setHeader("Authorization", "Bearer " + accessToken.getTokenValue());
return message;
});
return template;
});
}
}


Spring Cloud Vault предоставляет seamless интеграцию:
spring:
cloud:
vault:
host: ${VAULT_HOST}
port: 8200
scheme: https
authentication: APPROLE
app-role:
role-id: ${VAULT_ROLE_ID}
secret-id: ${VAULT_SECRET_ID}
kv:
backend: kv-v2
default-context: rabbitmq/production


Динамическое получение credentials:
@Configuration
public class VaultRabbitMQConfig {

@Bean
@VaultPropertySource(
value = "rabbitmq/production/credentials",
renewal = Renewal.ROTATE
)
public PropertySource<?> vaultPropertySource() {
return new MapPropertySource("vault-rabbitmq", Map.of());
}

@Bean
@RefreshScope
public ConnectionFactory connectionFactory(
@Value("${spring.rabbitmq.vault.username}") String username,
@Value("${spring.rabbitmq.vault.password}") String password) {

CachingConnectionFactory factory = new CachingConnectionFactory();
factory.setUsername(username);
factory.setPassword(password);

// Автоматическое обновление при rotation credentials
factory.setConnectionListeners(Collections.singletonList(connection -> {
log.info("RabbitMQ connection updated with new credentials");
}));

return factory;
}
}


mTLS с автоматической ротацией сертификатов
@Configuration
public class MTLSConfig {

@Bean
@Scheduled(fixedDelay = 3600000) // Каждый час
public void rotateMTLSCertificates() {
// Автоматическая ротация сертификатов mTLS
// Интеграция с cert-manager или аналогичным решением
}

@Bean
public SslContextFactory sslContextFactory(VaultTemplate vaultTemplate) {
return SslContextFactory.builder()
.keyManager(keyManagerFromVault(vaultTemplate))
.trustManager(trustManagerFromVault(vaultTemplate))
.protocol("TLSv1.3")
.build();
}

private KeyManager keyManagerFromVault(VaultTemplate vaultTemplate) {
// Получение клиентского сертификата из Vault
VaultResponse response = vaultTemplate.read("pki/issue/rabbitmq-client");

byte[] certificate = Base64.getDecoder()
.decode(response.getData().get("certificate").toString());
byte[] privateKey = Base64.getDecoder()
.decode(response.getData().get("private_key").toString());

return loadKeyManager(certificate, privateKey);
}
}


#Java #middle #RabbitMQ
👍3🔥1
Что выведет код?

public class Task050126 {
public static void main(String[] args) throws InterruptedException {
Counter c1 = new Counter();
Counter c2 = new Counter();

Thread t1 = new Thread(() -> c1.increment());
Thread t2 = new Thread(() -> c2.increment());

t1.start();
t2.start();

t1.join();
t2.join();

System.out.println(Counter.count);
}

static class Counter {
static int count = 0;
private final Object lock = new Object();

public void increment() {
for (int i = 0; i < 100000; i++) {
synchronized (lock) {
count++;
}
}
}
}
}


#Tasks
👍4
👍33
Вопрос с собеседований

Когда ForkJoinPool может ухудшить производительность? 🤓

Ответ:

Если задачи блокирующие (I/O, synchronized), ForkJoinPool теряет эффективность.

Он рассчитан на CPU-bound задачи. Блокировки приводят к starvation и снижению параллелизма, особенно при использовании parallelStream.


#собеседование
Please open Telegram to view this post
VIEW IN TELEGRAM
👍5
История IT-технологий сегодня — 06 января


ℹ️ Кто родился в этот день

Не нашел(


🌐 Знаковые события

2015 — открытие экзопланеты с самым высоким Индексом подобия Земле — Kepler-438 b.

#Biography #Birth_Date #Events #06Января
Please open Telegram to view this post
VIEW IN TELEGRAM
👍4
Раздел 6. Коллекции в Java

Глава 8. Дополнительные аспекты коллекций

Практика: В «Библиотеке» сделать коллекцию книг потокобезопасной (CopyOnWriteArrayList). Реализовать неизменяемый список популярных книг для чтения

Перед началом убедитесь, что проект готов, и вспомните ключевые концепции:
CopyOnWriteArrayList: Thread-safe версия List, где модификации создают копию массива (copy-on-write), а чтение — без locks. Идеально для read-heavy сценариев (много чтения, мало записи).
Преимущества: Безопасность без синхронизации, итераторы не fail-fast (не бросают ConcurrentModificationException при mod).
Недостатки: Высокий overhead на память и время для модификаций (копия всего списка), не подходит для write-heavy.

Неизменяемые коллекции: Collections.unmodifiableList делает List read-only — методы mod бросают UnsupportedOperationException. Полезно для constants или защиты данных.
Импорты: java.util.concurrent.CopyOnWriteArrayList, java.util.Collections.


Откройте проект

Запустите IDE, откройте LibraryProject. Проверьте List<Book> books, методы addBook, printAllBooks и т.д.
Импортируйте пакеты: В Library.java добавьте import java.util.concurrent.CopyOnWriteArrayList; и import java.util.Collections;. IDE поможет.
Планирование: Мы изменим поле books на CopyOnWriteArrayList, добавим метод для создания неизменяемого списка популярных книг и протестируем в multi-thread (опционально).

Замена List на CopyOnWriteArrayList для потокобезопасности
CopyOnWriteArrayList — отличный выбор для библиотеки, где чтение (поиск, вывод) частое, а запись (добавление книг) редкое. Это сделает коллекцию thread-safe без manual locks.


Измените поле books

В классе Library замените private List<Book> books = new ArrayList<>(); на private CopyOnWriteArrayList<Book> books = new CopyOnWriteArrayList<>();.
Это обеспечит thread-safety: при add/remove создается копия внутреннего массива, читатели видят snapshot.


Обновите конструктор Library


Если инициализация в конструкторе — обновите на CopyOnWriteArrayList.


Обновите метод addBook(Book book)

Используйте books.add(book); — это thread-safe (копия под капотом).
Добавьте проверку: if (book == null) return; или бросьте IllegalArgumentException.
Выведите сообщение о добавлении.


Обновите методы, использующие books

В printAllBooks() или findBookByTitle используйте for-each или Iterator — они работают на snapshot, безопасны при параллельных mod.
В removeBookByIndex(int index): books.remove(index); — thread-safe.


Реализация неизменяемого списка популярных книг

Неизменяемый список — это read-only view, полезный для "популярных" книг, которые не должны изменяться.

Добавьте поле для популярных книг:
В Library добавьте приватное поле List<Book> popularBooks = new ArrayList<>();.
В конструкторе или методе initPopularBooks() добавьте 3-5 статических книг (new Book("1984", "Orwell", 1949) и т.д.).

Сделайте список неизменяемым:

После заполнения popularBooks присвойте ему Collections.unmodifiableList(popularBooks);.
Теперь методы mod (add, remove) бросят UnsupportedOperationException.

Добавьте метод getPopularBooks():
Возвращайте unmodifiable список — public List<Book> getPopularBooks() { return popularBooks; }.
Это безопасно — внешний код не сможет изменить.

Добавьте метод printPopularBooks():
Переберите unmodifiable список и выведите детали книг.
Попробуйте в Main popularBooks.add(new Book(...)) — поймайте исключение.


#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList #Практика
👍4
Тестирование и отладка (объемно)

Базовое тестирование:
В Main добавьте книги, вызовите printAllBooks — всё как раньше.

Неизменяемость тест:
В Main получите getPopularBooks(), попробуйте add/remove — UnsupportedOperationException.
Переберите и выведите — чтение работает.

Отладка нюансов:
В CopyOnWriteArrayList добавьте breakpoint в add — увидите копию массива.
Тестируйте с большим размером (1000 элементов) — замерьте время (System.nanoTime()) для add в multi-thread.
Ловушки: CopyOnWriteArrayList slow для write-heavy — если много добавлений, используйте synchronized List.

Эксперименты (объемно и превышает текущий уровень экспертизы):

Замените на synchronizedList(new ArrayList<>()) — протестируйте thread-safety, но заметьте locks (медленнее для read-heavy).
Добавьте метод addPopularBook(Book book) — но сделайте так, чтобы он работал только до unmodifiable (или используйте builder для init).
Тестируйте с null — CopyOnWriteArrayList позволяет null элементы.
В multi-thread добавьте System.out в add/print — увидите interleaving, но без ошибок.
Измерьте память (Runtime.getRuntime().totalMemory()) перед/после многих add — увидите overhead копий в CopyOnWrite.


Работу данных коллекций мы будем проверять когда дойдем до многопоточности в java.


#Java #для_новичков #beginner #immutability #Collection #synchronizedList #ConcurrentHashMap #CopyOnWriteArrayList #Практика
👍3
Варианты ответа:
Anonymous Quiz
0%
1
41%
2
12%
null
47%
Исключение
👍2
Что выведет код?

import java.util.concurrent.ConcurrentHashMap;

public class Task060126 {
public static void main(String[] args) {
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();

Integer result = map.computeIfAbsent("a", k -> map.computeIfAbsent("a", k2 -> 2));

System.out.println(result);
}
}


#Tasks
👍3
Вопрос с собеседований

Почему hashCode() должен быть согласован с equals()? 🤓

Ответ:

Если equals возвращает true, hashCode обязан быть одинаковым.

Иначе HashMap/HashSet не смогут корректно находить элементы.

Нарушение контракта приводит к потерянным объектам и трудноуловимым багам.


#собеседование
Please open Telegram to view this post
VIEW IN TELEGRAM
👍7
История IT-технологий сегодня — 07 января


ℹ️ Кто родился в этот день

Иоганн Филипп Рейс (нем. Johann Philipp Reis; 7 января 1834, Гельнхаузен, Великое герцогство Гессен — 14 января 1874, Фридрихсдорф, Германия) — немецкий физик и изобретатель, первым в 1860 году сконструировавший электрический телефон, который в его честь сейчас называется телефоном Рейса. Впервые это изобретение было продемонстрировано публике 25 октября 1861 года.


🌐 Знаковые события

1610 — открытие четырёх крупнейших спутников Юпитера Галилео Галилеем.

1954 – Эксперимент Джорджтаун-IBM: Первая публичная демонстрация системы машинного перевода состоялась в Нью-Йорке в головном офисе IBM.


#Biography #Birth_Date #Events #07Января
Please open Telegram to view this post
VIEW IN TELEGRAM
👍4
Паттерны использования и гарантии доставки в RabbitMQ

RabbitMQ предоставляет гибкие модели взаимодействия, позволяющие реализовывать различные архитектурные паттерны. Понимание этих паттернов и связанных с ними гарантий доставки критически важно для построения надежных распределенных систем.


Work Queues: распределение нагрузки между потребителями

Work Queues (очереди задач) используются для распределения трудоемких задач между несколькими worker-процессами. Этот паттерн идеален для обработки фоновых задач, таких как генерация отчетов, обработка изображений или отправка email.

Ключевые характеристики:
Одна очередь, несколько потребителей
Каждое сообщение обрабатывается только одним потребителем
Балансировка нагрузки через настройку prefetch count

Реализация на Java с Spring AMQP

Producer:
@Service
public class TaskProducer {
private final RabbitTemplate rabbitTemplate;

@Value("${app.queues.task-queue}")
private String taskQueue;

public void sendTask(Task task) {
rabbitTemplate.convertAndSend(taskQueue, task, message -> {
// Установка приоритета задачи
message.getMessageProperties().setPriority(task.getPriority());
// Время жизни сообщения
message.getMessageProperties().setExpiration("3600000"); // 1 час
return message;
});
}
}


Consumer с ручными подтверждениями:
@Component
public class TaskWorker {

@RabbitListener(queues = "${app.queues.task-queue}")
public void processTask(Task task, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
try {
// Обработка задачи
boolean success = executeTask(task);

if (success) {
// Подтверждение успешной обработки
channel.basicAck(deliveryTag, false);
log.info("Task {} processed successfully", task.getId());
} else {
// Отказ без повторной очереди (перемещение в DLQ)
channel.basicNack(deliveryTag, false, false);
log.error("Task {} failed, moved to DLQ", task.getId());
}
} catch (Exception e) {
// При ошибке - повторная очередь
channel.basicNack(deliveryTag, false, true);
log.error("Error processing task {}, requeued", task.getId(), e);
}
}

private boolean executeTask(Task task) {
// Логика обработки задачи
return true;
}
}


Конфигурация для равномерного распределения:
@Configuration
public class WorkQueueConfig {

@Bean
public Queue taskQueue() {
return QueueBuilder.durable("tasks.queue")
.withArgument("x-max-priority", 10) // Поддержка приоритетов
.withArgument("x-dead-letter-exchange", "dlx.exchange")
.build();
}

@Bean
public SimpleRabbitListenerContainerFactory workerFactory(
ConnectionFactory connectionFactory) {

SimpleRabbitListenerContainerFactory factory =
new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory);

// Критичная настройка для равномерного распределения
factory.setPrefetchCount(1); // По одному сообщению на consumer
factory.setConcurrentConsumers(3); // Три параллельных worker'а
factory.setMaxConcurrentConsumers(10); // Автомасштабирование при нагрузке

// Ручные подтверждения для контроля
factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);

return factory;
}
}



#Java #middle #RabbitMQ
👍4
Publish/Subscribe: широковещательная рассылка

Publish/Subscribe (публикация/подписка) используется, когда одно сообщение должно быть доставлено множеству потребителей. Типичные сценарии: уведомления, аудит-логи, обновления кэша.

Архитектура:
Publisher отправляет сообщение в exchange
Exchange копирует сообщение во все привязанные очереди
Каждый consumer имеет свою собственную очередь

Реализация Fanout Exchange

Конфигурация обменника и очередей:

@Configuration
public class PubSubConfig {

@Bean
public FanoutExchange notificationsExchange() {
return new FanoutExchange("notifications.fanout", true, false);
}

@Bean
public Queue auditQueue() {
return new Queue("notifications.audit.queue", true);
}

@Bean
public Queue cacheQueue() {
return new Queue("notifications.cache.queue", true);
}

@Bean
public Queue analyticsQueue() {
return new Queue("notifications.analytics.queue", true);
}

@Bean
public Binding auditBinding() {
return BindingBuilder.bind(auditQueue())
.to(notificationsExchange());
}

@Bean
public Binding cacheBinding() {
return BindingBuilder.bind(cacheQueue())
.to(notificationsExchange());
}

@Bean
public Binding analyticsBinding() {
return BindingBuilder.bind(analyticsQueue())
.to(notificationsExchange());
}
}


Publisher:
@Service
public class NotificationPublisher {
private final RabbitTemplate rabbitTemplate;

public void publishNotification(Notification notification) {
// Отправка в fanout exchange - все подписчики получат копию
rabbitTemplate.convertAndSend(
"notifications.fanout",
"", // routing key игнорируется для fanout
notification
);
}
}


Consumer для аудита (пример одного из многих):
@Component
public class AuditNotificationConsumer {

@RabbitListener(queues = "notifications.audit.queue")
public void handleNotification(Notification notification) {
// Каждый consumer обрабатывает свою копию сообщения
auditService.logEvent(notification);
}
}



Routing/Topics: селективная подписка

Topic Exchange позволяет выполнять сложную маршрутизацию на основе routing key и patterns. Используется, когда потребители интересуются только определенными типами сообщений.

Pattern syntax:
* (звездочка) заменяет одно слово
# (решетка) заменяет ноль или более слов
Пример: stock.usd.* или stock.#

Реализация Topic Exchange

Конфигурация:
@Configuration
public class TopicRoutingConfig {

@Bean
public TopicExchange stockExchange() {
return new TopicExchange("stock.topic", true, false);
}

@Bean
public Queue usdStockQueue() {
return new Queue("stock.usd.queue", true);
}

@Bean
public Queue eurStockQueue() {
return new Queue("stock.eur.queue", true);
}

@Bean
public Queue allStockQueue() {
return new Queue("stock.all.queue", true);
}

@Bean
public Binding usdBinding() {
// Будет получать: stock.usd.nyse, stock.usd.nasdaq
return BindingBuilder.bind(usdStockQueue())
.to(stockExchange())
.with("stock.usd.*");
}

@Bean
public Binding eurBinding() {
// Будет получать: stock.eur.lse, stock.eur.euronext
return BindingBuilder.bind(eurStockQueue())
.to(stockExchange())
.with("stock.eur.*");
}

@Bean
public Binding allStocksBinding() {
// Будет получать все сообщения о stock
return BindingBuilder.bind(allStockQueue())
.to(stockExchange())
.with("stock.#");
}
}


#Java #middle #RabbitMQ
👍4
Publisher с различными routing keys:
@Service
public class StockPublisher {
private final RabbitTemplate rabbitTemplate;

public void publishUSDPrice(String exchange, BigDecimal price) {
String routingKey = "stock.usd." + exchange.toLowerCase();
StockUpdate update = new StockUpdate("USD", exchange, price);

rabbitTemplate.convertAndSend(
"stock.topic",
routingKey,
update
);
}

public void publishEURPrice(String exchange, BigDecimal price) {
String routingKey = "stock.eur." + exchange.toLowerCase();
StockUpdate update = new StockUpdate("EUR", exchange, price);

rabbitTemplate.convertAndSend(
"stock.topic",
routingKey,
update
);
}
}


RPC over AMQP: синхронные вызовы через асинхронный транспорт

Remote Procedure Call (RPC) поверх AMQP позволяет выполнять синхронные запросы-ответы через ассинхронную систему сообщений. Используется когда нужен немедленный ответ, но нельзя установить прямое соединение.

Механизм работы:
Клиент отправляет запрос с уникальным correlationId
Сервер обрабатывает запрос и отправляет ответ в очередь ответов
Клиент ожидает ответ с matching correlationId

Реализация RPC

Клиентская сторона:
@Service
public class CalculatorRpcClient {
private final RabbitTemplate rabbitTemplate;

public CalculatorRpcClient(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;

// Настройка reply listener
this.rabbitTemplate.setReplyTimeout(30000);
this.rabbitTemplate.setUseDirectReplyToContainer(false);
}

public BigDecimal calculate(CalculationRequest request) {
// Создание correlation ID
String correlationId = UUID.randomUUID().toString();

// Отправка запроса и ожидание ответа
CalculationResponse response = (CalculationResponse)
rabbitTemplate.convertSendAndReceive(
"rpc.exchange",
"calculator.rpc",
request,
message -> {
message.getMessageProperties()
.setCorrelationId(correlationId);
message.getMessageProperties()
.setReplyTo("amq.rabbitmq.reply-to"); // Временная очередь
return message;
}
);

if (response == null) {
throw new RuntimeException("RPC timeout");
}

return response.getResult();
}
}


Серверная сторона:
@Component
public class CalculatorRpcServer {

@RabbitListener(
bindings = @QueueBinding(
value = @Queue(value = "rpc.calculator.queue", durable = "true"),
exchange = @Exchange(value = "rpc.exchange", type = ExchangeTypes.DIRECT),
key = "calculator.rpc"
)
)
public CalculationResponse calculate(CalculationRequest request,
Message message) {

String correlationId = message.getMessageProperties()
.getCorrelationId();
String replyTo = message.getMessageProperties()
.getReplyTo();

log.debug("Processing RPC request {} from {}",
correlationId, replyTo);

// Выполнение расчета
BigDecimal result = performCalculation(request);

// Возврат результата с сохранением correlationId
CalculationResponse response = new CalculationResponse(result);

// CorrelationId автоматически копируется в ответ
return response;
}

private BigDecimal performCalculation(CalculationRequest request) {
// Логика расчета
return BigDecimal.ZERO;
}
}



#Java #middle #RabbitMQ
👍4
Гарантии доставки

At-most-once: риск потери сообщений

Механизм: Сообщение отправляется без подтверждений, сразу удаляется из очереди.

Использовать когда:

Потеря сообщений допустима
Высокая производительность критична
Система имеет механизмы компенсации

// Producer без подтверждений
@Bean
public RabbitTemplate atMostOnceTemplate(ConnectionFactory cf) {
RabbitTemplate template = new RabbitTemplate(cf);
template.setMandatory(false);

// Отключение publisher confirms
((CachingConnectionFactory) cf)
.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.NONE);

return template;
}

// Consumer с auto-ack
@RabbitListener(queues = "at.most.once.queue")
public void handleAutoAck(Message message) {
// Сообщение уже удалено из очереди при доставке
// При падении здесь - сообщение потеряно
processMessage(message);
}


At-least-once: риск дублирования

Механизм: Гарантирует доставку минимум один раз, но возможны дубликаты.

Реализация с ручными подтверждениями:
@Component
public class AtLeastOnceConsumer {

@RabbitListener(queues = "orders.queue")
public void processOrder(Order order, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {

try {
// 1. Обработка заказа
processOrderSafely(order);

// 2. Подтверждение после успешной обработки
channel.basicAck(tag, false);

log.info("Order {} processed and acknowledged", order.getId());

} catch (Exception e) {
// 3. При ошибке - повторная очередь
channel.basicNack(tag, false, true);
log.error("Order processing failed, requeued", e);

// Идемпотентность критична!
// При повторной обработке должна быть сохранена консистентность
}
}

private void processOrderSafely(Order order) {
// Идемпотентная обработка
if (orderRepository.existsById(order.getId())) {
log.warn("Duplicate order {}, skipping", order.getId());
return;
}

orderRepository.save(order);
inventoryService.reserveItems(order);
}
}



#Java #middle #RabbitMQ
👍4
Exactly-once: сложная реализация

Реальность: RabbitMQ не предоставляет exactly-once гарантий на уровне протокола. Достигается комбинацией механизмов.

Паттерн для pseudo-exactly-once:

@Service
@Transactional
public class ExactlyOnceProcessor {

private final RabbitTemplate rabbitTemplate;
private final OrderRepository orderRepository;

public void processWithExactlyOnceSemantics(OrderMessage message) {
// 1. Проверка идемпотентности
if (processedMessageRepository.existsById(message.getId())) {
return; // Уже обработано
}

// 2. Сохранение в outbox паттерн
OrderOutbox outbox = new OrderOutbox();
outbox.setMessageId(message.getId());
outbox.setOrderData(message.getPayload());
outbox.setStatus(OutboxStatus.PROCESSING);

orderRepository.save(outbox);

// 3. Бизнес-логика в транзакции
Order order = createOrder(message);
orderRepository.save(order);

// 4. Обновление статуса в той же транзакции
outbox.setStatus(OutboxStatus.PROCESSED);
orderRepository.save(outbox);

// 5. Отправка подтверждения (после коммита транзакции)
// RabbitTemplate работает после завершения транзакции
rabbitTemplate.convertAndSend(
"order.processed.exchange",
"order.processed",
new OrderProcessedEvent(order.getId())
);

// 6. Идемпотентность на стороне получателя подтверждения
}

@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void afterCommit(OrderProcessedEvent event) {
// Только после успешного коммита
rabbitTemplate.convertAndSend(
"order.completed.queue",
event
);
}
}


Механизмы надежности


Publisher Confirms

Механизм: Брокер подтверждает получение сообщения.

Реализация:

@Configuration
public class PublisherConfirmConfig {

@Bean
public CachingConnectionFactory confirmedConnectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory();

// Включение publisher confirms
factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
factory.setPublisherReturns(true);

return factory;
}

@Bean
public RabbitTemplate confirmedRabbitTemplate(
CachingConnectionFactory connectionFactory) {

RabbitTemplate template = new RabbitTemplate(connectionFactory);

// Callback для подтверждений
template.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
metrics.increment("publisher.confirms.success");
} else {
metrics.increment("publisher.confirms.failure");
log.error("Message not confirmed: {}", cause);

// Логика повторной отправки
if (correlationData != null) {
retryService.scheduleRetry(correlationData.getId());
}
}
});

// Callback для возвращенных сообщений
template.setReturnsCallback(returned -> {
log.error("Message returned: {}", returned.toString());
returnedMessageService.handleReturned(returned);
});

return template;
}
}



#Java #middle #RabbitMQ
👍4
Consumer Acknowledgements

Типы подтверждений:
public class AcknowledgementExamples {

// 1. Автоматическое подтверждение (ненадежно)
@RabbitListener(queues = "auto.ack.queue")
public void autoAck(Message message) {
// Подтверждение происходит автоматически при возврате метода
}

// 2. Ручное подтверждение (рекомендуется)
@RabbitListener(queues = "manual.ack.queue")
public void manualAck(Order order, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {

try {
process(order);
channel.basicAck(tag, false); // Положительное подтверждение
} catch (BusinessException e) {
// Не requeue для бизнес-ошибок
channel.basicNack(tag, false, false);
} catch (TechnicalException e) {
// Requeue для технических ошибок
channel.basicNack(tag, false, true);
}
}

// 3. Подтверждение с транзакцией
@Transactional
@RabbitListener(queues = "transactional.queue")
public void transactionalAck(Payment payment, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {

// Сохранение в БД
paymentRepository.save(payment);

// Подтверждение произойдет только после коммита транзакции
// При откате транзакции - сообщение вернется в очередь
}
}


Transactional Channels

Использование транзакций:

@Service
public class TransactionalMessageService {

private final RabbitTemplate rabbitTemplate;

@Transactional
public void processInTransaction(Order order) {
// 1. Сохранение в базу данных
orderRepository.save(order);

// 2. Отправка сообщения в той же транзакции
rabbitTemplate.execute(channel -> {
// Начало транзакции RabbitMQ
channel.txSelect();

try {
// Отправка сообщения
channel.basicPublish(
"orders.exchange",
"order.created",
null,
serialize(order)
);

// Коммит транзакции RabbitMQ
channel.txCommit();

return null;
} catch (Exception e) {
// Откат транзакции RabbitMQ
channel.txRollback();
throw new RuntimeException("Transaction failed", e);
}
});

// 3. Если произойдет исключение здесь,
// откатятся и БД и сообщение RabbitMQ
inventoryService.updateStock(order);
}
}



#Java #middle #RabbitMQ
👍4