🌱 Spring AMQP et Spring Boot
Spring AMQP offre une intégration élégante de RabbitMQ dans l'écosystème Spring, transformant la complexité de la messagerie AMQP en annotations simples et configuration déclarative. Cette intégration permet aux développeurs Java de tirer parti de la puissance de RabbitMQ sans quitter l'univers familier de Spring.
🧠 Philosophie d'Intégration Spring
Convention over Configuration
Spring AMQP applique la philosophie "Convention over Configuration" à la messagerie. Plutôt que de configurer manuellement chaque aspect de RabbitMQ, Spring fournit des defaults intelligents et des conventions qui réduisent drastiquement le boilerplate.
Auto-configuration magique : Spring Boot détecte automatiquement RabbitMQ et configure les beans nécessaires. Cette magie permet de démarrer rapidement tout en gardant la possibilité de customiser finement.
Abstraction sans perte de contrôle : L'abstraction Spring n'empêche pas d'accéder aux fonctionnalités avancées de RabbitMQ quand nécessaire. C'est une couche de commodité, pas une limitation.
🚀 Setup et Configuration
Démarrage Rapide
<!-- pom.xml -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
# application.yml - Configuration minimaliste
spring:
rabbitmq:
host: localhost
port: 5672
username: admin
password: secret
virtual-host: /myapp
@SpringBootApplication
@EnableRabbit // Active le support RabbitMQ
public class MessagingApplication {
public static void main(String[] args) {
SpringApplication.run(MessagingApplication.class, args);
}
}
Configuration Avancée
@Configuration
@EnableRabbitMQ
public class RabbitMQConfig {
// Connexion avec pool et retry
@Bean
public CachingConnectionFactory connectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory();
factory.setHost("rabbitmq-cluster");
factory.setPort(5672);
factory.setUsername("app-user");
factory.setPassword("secure-password");
factory.setVirtualHost("/production");
// Pool de connexions
factory.setChannelCacheSize(25);
factory.setConnectionCacheSize(3);
// Retry et recover
factory.getRabbitConnectionFactory().setAutomaticRecoveryEnabled(true);
factory.getRabbitConnectionFactory().setNetworkRecoveryInterval(5000);
factory.getRabbitConnectionFactory().setRequestedHeartbeat(60);
return factory;
}
// Template pour publishing
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
template.setMandatory(true); // Fail si non-routable
template.setConfirmCallback(this::publishConfirm);
template.setReturnCallback(this::publishReturn);
template.setMessageConverter(new Jackson2JsonMessageConverter());
// Retry template pour résilience
RetryTemplate retryTemplate = new RetryTemplate();
retryTemplate.setRetryPolicy(
new SimpleRetryPolicy(3, Map.of(AmqpException.class, true))
);
template.setRetryTemplate(retryTemplate);
return template;
}
// Listener container avec tuning
@Bean
public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory() {
SimpleRabbitListenerContainerFactory factory =
new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory());
factory.setConcurrentConsumers(3);
factory.setMaxConcurrentConsumers(10);
factory.setPrefetchCount(5); // Fair dispatch
factory.setDefaultRequeueRejected(false); // DLQ
factory.setErrorHandler(new CustomErrorHandler());
return factory;
}
private void publishConfirm(CorrelationData correlationData, boolean ack, String cause) {
if (ack) {
log.info("Message confirmed: {}", correlationData.getId());
} else {
log.error("Message rejected: {} - {}", correlationData.getId(), cause);
}
}
private void publishReturn(Message message, int replyCode, String replyText,
String exchange, String routingKey) {
log.error("Message returned: {} - {} -> {}:{}",
replyText, message, exchange, routingKey);
}
}
🏗️ Déclaration d'Infrastructure
Approach Déclarative
@Configuration
public class RabbitInfrastructureConfig {
// === EXCHANGES ===
@Bean
public TopicExchange userEventsExchange() {
return ExchangeBuilder
.topicExchange("user.events")
.durable(true)
.build();
}
@Bean
public DirectExchange userCommandsExchange() {
return ExchangeBuilder
.directExchange("user.commands")
.durable(true)
.build();
}
@Bean
public FanoutExchange notificationsExchange() {
return ExchangeBuilder
.fanoutExchange("notifications")
.durable(true)
.build();
}
// Dead Letter Exchange
@Bean
public DirectExchange deadLetterExchange() {
return ExchangeBuilder
.directExchange("dlx")
.durable(true)
.build();
}
// === QUEUES ===
@Bean
public Queue userRegistrationQueue() {
return QueueBuilder
.durable("user.registration")
.withArgument("x-dead-letter-exchange", "dlx")
.withArgument("x-dead-letter-routing-key", "user.registration.failed")
.withArgument("x-message-ttl", 1800000) // 30 minutes
.build();
}
@Bean
public Queue emailQueue() {
return QueueBuilder
.durable("notifications.email")
.withArgument("x-max-priority", 10)
.withArgument("x-dead-letter-exchange", "dlx")
.build();
}
@Bean
public Queue smsQueue() {
return QueueBuilder
.durable("notifications.sms")
.withArgument("x-max-priority", 10)
.withArgument("x-dead-letter-exchange", "dlx")
.build();
}
@Bean
public Queue deadLetterQueue() {
return QueueBuilder
.durable("dead.letters")
.build();
}
// === BINDINGS ===
@Bean
public Binding userRegistrationBinding() {
return BindingBuilder
.bind(userRegistrationQueue())
.to(userCommandsExchange())
.with("user.register");
}
@Bean
public Binding userCreatedEventBinding() {
return BindingBuilder
.bind(emailQueue())
.to(userEventsExchange())
.with("user.created.*");
}
@Bean
public Binding emailNotificationBinding() {
return BindingBuilder
.bind(emailQueue())
.to(notificationsExchange());
}
@Bean
public Binding smsNotificationBinding() {
return BindingBuilder
.bind(smsQueue())
.to(notificationsExchange());
}
@Bean
public Binding deadLetterBinding() {
return BindingBuilder
.bind(deadLetterQueue())
.to(deadLetterExchange())
.with("#"); // Tous les messages DLX
}
}
🎯 Programming Model
Producers Simplifiés
@Service
public class UserEventPublisher {
private final RabbitTemplate rabbitTemplate;
private final ObjectMapper objectMapper;
public UserEventPublisher(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
this.objectMapper = new ObjectMapper();
}
// Publishing simple
public void publishUserCreated(User user) {
UserCreatedEvent event = new UserCreatedEvent(
user.getId(),
user.getEmail(),
Instant.now()
);
rabbitTemplate.convertAndSend(
"user.events",
"user.created.new",
event
);
}
// Publishing avec propriétés
public void publishPriorityEvent(Object event, int priority) {
rabbitTemplate.convertAndSend(
"user.events",
"user.priority.action",
event,
message -> {
message.getMessageProperties().setPriority(priority);
message.getMessageProperties().setExpiration("300000"); // 5min TTL
message.getMessageProperties().setHeaders(Map.of(
"source", "user-service",
"timestamp", System.currentTimeMillis()
));
return message;
}
);
}
// Publishing avec confirmation
@Retryable(value = AmqpException.class, maxAttempts = 3)
public boolean publishWithConfirmation(String exchange, String routingKey, Object message) {
CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.convertAndSend(exchange, routingKey, message, correlationData);
// Attendre confirmation (pour cas critiques uniquement)
try {
CorrelationData.Confirm confirm = correlationData.getFuture().get(5, TimeUnit.SECONDS);
return confirm.isAck();
} catch (Exception e) {
log.error("Failed to get confirmation", e);
return false;
}
}
}
Consumers avec Annotations
@Component
public class UserEventHandlers {
private final EmailService emailService;
private final AnalyticsService analyticsService;
// Consumer simple
@RabbitListener(queues = "user.registration")
public void handleUserRegistration(UserRegistrationCommand command) {
log.info("Processing user registration: {}", command.getUserId());
// Traitement métier
User user = userService.createUser(command);
// Publier événement de succès
eventPublisher.publishUserCreated(user);
}
// Consumer avec gestion d'erreurs
@RabbitListener(queues = "notifications.email")
public void handleEmailNotification(
@Payload EmailNotification notification,
@Header Map<String, Object> headers,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag
) {
try {
emailService.sendEmail(notification);
channel.basicAck(deliveryTag, false);
log.info("Email sent successfully: {}", notification.getTo());
} catch (RetryableEmailException e) {
// Retry avec exponential backoff
int retryCount = (Integer) headers.getOrDefault("x-retry-count", 0);
if (retryCount < 3) {
republishWithDelay(notification, retryCount + 1);
channel.basicAck(deliveryTag, false);
} else {
channel.basicReject(deliveryTag, false); // Vers DLQ
}
} catch (PermanentEmailException e) {
log.error("Permanent email failure: {}", e.getMessage());
channel.basicReject(deliveryTag, false); // Direct vers DLQ
}
}
// Consumer conditionnel avec SpEL
@RabbitListener(
bindings = @QueueBinding(
value = @Queue(name = "vip.notifications", durable = "true"),
exchange = @Exchange(name = "user.events", type = "topic"),
key = "user.*.premium"
),
condition = "#{environment.getProperty('features.vip-notifications') == 'true'}"
)
public void handleVipNotifications(UserEvent event) {
log.info("VIP notification for premium user: {}", event.getUserId());
vipNotificationService.sendPremiumNotification(event);
}
// Dead Letter Queue handler
@RabbitListener(queues = "dead.letters")
public void handleFailedMessages(@Payload String failedMessage,
@Header Map<String, Object> headers) {
log.error("Processing failed message: {}", failedMessage);
// Analyser la cause d'échec
String originalQueue = (String) headers.get("x-first-death-queue");
String reason = (String) headers.get("x-first-death-reason");
// Alerter les opérations
alertingService.sendAlert("DLQ Message", Map.of(
"original_queue", originalQueue,
"reason", reason,
"message", failedMessage
));
// Optionnel : tentative de récupération
if (canRecover(reason)) {
recoverMessage(failedMessage, originalQueue);
}
}
}
🔧 Configuration Avancée
Profiles et Environnements
@Configuration
@Profile("development")
public class DevelopmentRabbitConfig {
@Bean
@Primary
public ConnectionFactory devConnectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory("localhost");
factory.setUsername("dev");
factory.setPassword("dev123");
// Config relaxe pour dev
factory.setChannelCacheSize(10);
factory.setPublisherConfirms(false); // Pas de confirms en dev
return factory;
}
}
@Configuration
@Profile("production")
public class ProductionRabbitConfig {
@Bean
@Primary
public ConnectionFactory prodConnectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory();
// Configuration cluster
factory.setAddresses("rabbit1:5672,rabbit2:5672,rabbit3:5672");
factory.setUsername(environment.getProperty("RABBITMQ_USER"));
factory.setPassword(environment.getProperty("RABBITMQ_PASSWORD"));
// Optimisations production
factory.setChannelCacheSize(50);
factory.setConnectionCacheSize(5);
factory.setPublisherConfirms(true);
factory.setPublisherReturns(true);
// SSL si nécessaire
if (environment.getProperty("RABBITMQ_SSL", Boolean.class, false)) {
factory.getRabbitConnectionFactory().useSslProtocol();
}
return factory;
}
// Admin pour management programmatique
@Bean
public AmqpAdmin amqpAdmin(ConnectionFactory connectionFactory) {
return new RabbitAdmin(connectionFactory);
}
}
Health Checks et Actuator
@Component
public class RabbitMQHealthIndicator implements HealthIndicator {
private final RabbitTemplate rabbitTemplate;
@Override
public Health health() {
try {
// Test de connectivité
rabbitTemplate.execute(channel -> {
channel.queueDeclarePassive("health.check.queue");
return null;
});
// Vérifier les métriques critiques
Map<String, Object> metrics = getSystemMetrics();
if ((Integer) metrics.get("memory_usage_percent") > 80) {
return Health.down()
.withDetail("reason", "High memory usage")
.withDetails(metrics)
.build();
}
return Health.up()
.withDetails(metrics)
.build();
} catch (Exception e) {
return Health.down()
.withDetail("error", e.getMessage())
.build();
}
}
private Map<String, Object> getSystemMetrics() {
// Intégration avec Management API
return Map.of(
"connections", getConnectionCount(),
"channels", getChannelCount(),
"queues", getQueueCount(),
"memory_usage_percent", getMemoryUsagePercent()
);
}
}
🎯 Patterns d'Implémentation
Event-Driven Microservices
// Service utilisateur avec événements
@Service
@Transactional
public class UserService {
private final UserRepository userRepository;
private final UserEventPublisher eventPublisher;
public User createUser(CreateUserCommand command) {
// 1. Validation métier
validateUserCreation(command);
// 2. Persistence (transaction)
User user = new User(command);
user = userRepository.save(user);
// 3. Publication événement (après commit)
eventPublisher.publishUserCreated(user);
return user;
}
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void handleUserCreated(UserCreatedEvent event) {
// Publier après commit de transaction
rabbitTemplate.convertAndSend(
"user.events",
"user.created.new",
event
);
}
}
// Service de notification réactif
@Component
public class NotificationService {
@RabbitListener(queues = "notifications.email")
@Transactional
public void handleUserCreated(UserCreatedEvent event) {
// Email de bienvenue
EmailTemplate template = templateService.getWelcomeTemplate();
Email email = template.render(event.getUser());
emailService.send(email);
// Analytics
analyticsService.trackEvent("user.welcome.email.sent", Map.of(
"userId", event.getUserId(),
"timestamp", event.getTimestamp()
));
}
@RabbitListener(queues = "notifications.sms")
public void handlePriorityNotification(
@Payload PriorityNotification notification,
@Header("priority") Integer priority
) {
if (priority != null && priority > 5) {
smsService.sendImmediate(notification);
} else {
smsService.sendBatch(notification);
}
}
}
Saga Pattern Implementation
// Orchestration de saga avec Spring
@Component
public class OrderSagaOrchestrator {
@RabbitListener(queues = "order.saga.start")
public void startOrderSaga(OrderCreatedEvent event) {
SagaTransaction saga = new SagaTransaction(event.getOrderId());
try {
// Étape 1: Réserver inventaire
saga.addStep("inventory", () ->
reserveInventory(event.getOrderId(), event.getItems())
);
// Étape 2: Traiter paiement
saga.addStep("payment", () ->
processPayment(event.getOrderId(), event.getAmount())
);
// Étape 3: Créer expédition
saga.addStep("shipping", () ->
createShipment(event.getOrderId(), event.getAddress())
);
// Exécuter la saga
saga.execute();
} catch (SagaException e) {
// Compensation automatique
saga.compensate();
publishOrderFailed(event.getOrderId(), e.getMessage());
}
}
private void reserveInventory(String orderId, List<OrderItem> items) {
ReserveInventoryCommand command = new ReserveInventoryCommand(orderId, items);
// Publication synchrone avec timeout
Object result = rabbitTemplate.convertSendAndReceive(
"inventory.commands",
"inventory.reserve",
command,
message -> {
message.getMessageProperties().setReplyTo("saga.inventory.reply");
message.getMessageProperties().setExpiration("30000"); // 30s timeout
return message;
}
);
if (result == null) {
throw new SagaException("Inventory reservation timeout");
}
InventoryReservationResult reservation = (InventoryReservationResult) result;
if (!reservation.isSuccess()) {
throw new SagaException("Inventory reservation failed: " + reservation.getReason());
}
}
}
🔍 Testing et Development
Test d'Intégration
@SpringBootTest
@TestMethodOrder(OrderAnnotation.class)
class RabbitMQIntegrationTest {
@Autowired
private RabbitTemplate rabbitTemplate;
@Autowired
private UserService userService;
@Test
@Order(1)
void shouldPublishUserCreatedEvent() throws InterruptedException {
// Arrange
CreateUserCommand command = new CreateUserCommand("john@example.com", "John Doe");
CountDownLatch latch = new CountDownLatch(1);
List<UserCreatedEvent> receivedEvents = new ArrayList<>();
// Setup test listener
setupTestListener(latch, receivedEvents);
// Act
userService.createUser(command);
// Assert
assertTrue(latch.await(5, TimeUnit.SECONDS));
assertEquals(1, receivedEvents.size());
assertEquals("john@example.com", receivedEvents.get(0).getEmail());
}
private void setupTestListener(CountDownLatch latch, List<UserCreatedEvent> events) {
rabbitTemplate.setReceiveTimeout(1000);
// Listener temporaire pour test
SimpleMessageListenerContainer container =
new SimpleMessageListenerContainer(connectionFactory);
container.setQueueNames("test.user.events");
container.setMessageListener(new MessageListener() {
@Override
public void onMessage(Message message) {
try {
UserCreatedEvent event = objectMapper.readValue(
message.getBody(), UserCreatedEvent.class
);
events.add(event);
latch.countDown();
} catch (Exception e) {
log.error("Test listener error", e);
}
}
});
container.start();
}
}
TestContainers Integration
@SpringBootTest
@Testcontainers
class RabbitMQTestContainersTest {
@Container
static RabbitMQContainer rabbitmq = new RabbitMQContainer("rabbitmq:3.12-management")
.withAdminPassword("admin123")
.withPluginsEnabled("rabbitmq_delayed_message_exchange");
@DynamicPropertySource
static void configureProperties(DynamicPropertyRegistry registry) {
registry.add("spring.rabbitmq.host", rabbitmq::getHost);
registry.add("spring.rabbitmq.port", rabbitmq::getAmqpPort);
registry.add("spring.rabbitmq.username", () -> "admin");
registry.add("spring.rabbitmq.password", () -> "admin123");
}
@Test
void shouldHandleMessageFlow() {
// Test complet avec vraie instance RabbitMQ
// Infrastructure créée automatiquement
// Nettoyage automatique après test
}
}
📊 Monitoring et Métriques
Intégration Actuator
@Component
public class RabbitMQMetrics {
@EventListener
public void handlePublishMetrics(RabbitPublishEvent event) {
Metrics.counter("rabbitmq.messages.published",
Tags.of(
"exchange", event.getExchange(),
"routing_key", event.getRoutingKey(),
"status", event.isSuccess() ? "success" : "failure"
)
).increment();
}
@EventListener
public void handleConsumeMetrics(RabbitConsumeEvent event) {
Timer.Sample sample = Timer.start(Metrics.globalRegistry);
// Mesurer temps de traitement
sample.stop(Timer.builder("rabbitmq.messages.processing.duration")
.tag("queue", event.getQueueName())
.tag("consumer", event.getConsumerName())
.register(Metrics.globalRegistry));
}
// Métriques custom
@Scheduled(fixedRate = 30000)
public void collectQueueMetrics() {
amqpAdmin.getQueueInfo("critical.orders").ifPresent(info -> {
Metrics.gauge("rabbitmq.queue.depth",
Tags.of("queue", "critical.orders"),
info.getMessageCount()
);
});
}
}
🚀 Production Patterns
Graceful Shutdown
@Component
public class GracefulShutdownManager {
private final SimpleRabbitListenerContainerFactory containerFactory;
private final List<SimpleMessageListenerContainer> containers = new ArrayList<>();
@EventListener
public void handleShutdown(ContextClosedEvent event) {
log.info("Initiating graceful shutdown...");
// 1. Arrêter d'accepter nouveaux messages
containers.forEach(container -> {
container.stop();
log.info("Stopped container: {}", container.getQueueNames());
});
// 2. Attendre fin des traitements en cours
containers.forEach(this::waitForCompletion);
// 3. Fermer connexions proprement
connectionFactory.destroy();
log.info("Graceful shutdown completed");
}
private void waitForCompletion(SimpleMessageListenerContainer container) {
try {
// Attendre max 30 secondes
for (int i = 0; i < 30; i++) {
if (container.getActiveConsumerCount() == 0) {
break;
}
Thread.sleep(1000);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
✅ Bonnes Pratiques Spring AMQP
1. Configuration Externalisée
# application.yml
spring:
rabbitmq:
addresses: ${RABBITMQ_ADDRESSES:localhost:5672}
username: ${RABBITMQ_USERNAME:guest}
password: ${RABBITMQ_PASSWORD:guest}
virtual-host: ${RABBITMQ_VHOST:/}
connection-timeout: 30s
# Publisher settings
publisher-confirms: true
publisher-returns: true
# Consumer settings
listener:
simple:
concurrency: ${RABBITMQ_CONSUMER_CONCURRENCY:3}
max-concurrency: ${RABBITMQ_CONSUMER_MAX_CONCURRENCY:10}
prefetch: ${RABBITMQ_CONSUMER_PREFETCH:5}
default-requeue-rejected: false
# Template settings
template:
mandatory: true
receive-timeout: 5s
reply-timeout: 10s
2. Error Handling Strategy
@Configuration
public class ErrorHandlingConfig {
@Bean
public RabbitListenerErrorHandler customErrorHandler() {
return (amqpMessage, message, exception) -> {
log.error("Listener error: {}", exception.getMessage(), exception);
// Custom recovery logic
if (exception instanceof MessageConversionException) {
// Log et ignorer les messages mal formés
deadLetterService.logBadMessage(amqpMessage, exception);
return null; // Ack le message
}
// Re-throw pour retry par défaut
throw exception;
};
}
// Custom retry template
@Bean
public RetryTemplate rabbitRetryTemplate() {
RetryTemplate template = new RetryTemplate();
// Policy: 3 tentatives avec backoff exponentiel
ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
backOffPolicy.setInitialInterval(1000);
backOffPolicy.setMultiplier(2.0);
backOffPolicy.setMaxInterval(10000);
SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(3, Map.of(
AmqpRejectAndDontRequeueException.class, false, // Pas de retry
MessageConversionException.class, false, // Pas de retry
Exception.class, true // Retry par défaut
));
template.setBackOffPolicy(backOffPolicy);
template.setRetryPolicy(retryPolicy);
return template;
}
}
🎯 Cas d'Usage Concrets
E-Commerce avec Spring Boot
// Architecture complète e-commerce
@SpringBootApplication
@EnableScheduling
@EnableRabbit
public class ECommerceApplication {
public static void main(String[] args) {
SpringApplication.run(ECommerceApplication.class, args);
}
// Order processing workflow
@Component
public class OrderWorkflow {
@RabbitListener(queues = "orders.new")
@Transactional
public void processNewOrder(OrderCreatedEvent event) {
Order order = orderService.findById(event.getOrderId());
// Workflow steps via messaging
publishInventoryCheck(order);
publishPaymentProcessing(order);
publishShippingPreparation(order);
}
@RabbitListener(queues = "inventory.responses")
public void handleInventoryResponse(InventoryCheckResult result) {
if (result.isAvailable()) {
orderService.markInventoryReserved(result.getOrderId());
} else {
orderService.cancelOrder(result.getOrderId(), "Stock indisponible");
}
}
}
}
Spring AMQP transforme RabbitMQ en une extension naturelle de votre application Spring, permettant de construire des systèmes distribués robustes avec la simplicité légendaire de Spring Boot.