📦 Propriétés des Messages

Les propriétés des messages dans RabbitMQ sont des métadonnées essentielles qui contrôlent le comportement, le routage et la livraison des messages. Comprendre et maîtriser ces propriétés est crucial pour construire des systèmes de messagerie robustes et performants.

🏗️ Structure d'un Message AMQP

Rendu du diagramme en cours...

Un message AMQP se compose de propriétés et d'un corps de message. Les propriétés se divisent en deux catégories :

  • Basic Properties : Propriétés standard définies par AMQP
  • Headers : Propriétés personnalisées définies par l'application

📋 Basic Properties Essentielles

Content Type et Encoding

// Publication avec content type
await channel.publish(
    'events',
    'user.created',
    Buffer.from(JSON.stringify(userData)),
    {
        contentType: 'application/json',
        contentEncoding: 'utf-8'
    }
);

// Support de différents formats
const messageFormats = {
    json: {
        contentType: 'application/json',
        serialize: JSON.stringify,
        deserialize: JSON.parse
    },
    xml: {
        contentType: 'application/xml', 
        serialize: obj => xmlBuilder.build(obj),
        deserialize: xml => xmlParser.parse(xml)
    },
    protobuf: {
        contentType: 'application/x-protobuf',
        serialize: obj => ProtoBuf.encode(obj),
        deserialize: bytes => ProtoBuf.decode(bytes)
    }
};

Delivery Mode (Persistance)

import pika

# Message persistant (survit au redémarrage)
channel.basic_publish(
    exchange='orders',
    routing_key='process',
    body=json.dumps(order),
    properties=pika.BasicProperties(
        delivery_mode=2  # 2 = persistant, 1 = non-persistant
    )
)

# Configuration complète pour production
def publish_critical_message(message_data):
    properties = pika.BasicProperties(
        delivery_mode=2,           # Persistant
        priority=9,                # Haute priorité
        timestamp=int(time.time()), # Timestamp
        message_id=str(uuid4()),   # ID unique
        user_id='system',          # Utilisateur émetteur
        app_id='order-service'     # Application émettrice
    )
    
    channel.basic_publish(
        exchange='critical.orders',
        routing_key='urgent.process',
        body=json.dumps(message_data),
        properties=properties,
        mandatory=True  # Échec si non routable
    )

TTL (Time To Live)

# TTL au niveau message
channel.basic_publish(
    exchange='notifications', 
    routing_key='email',
    body='Welcome email',
    properties=pika.BasicProperties(
        expiration='300000'  # 5 minutes en millisecondes
    )
)

# TTL au niveau queue (plus efficace)
channel.queue_declare(
    queue='temporary_notifications',
    durable=True,
    arguments={
        'x-message-ttl': 600000,  # 10 minutes pour tous les messages
        'x-expires': 1800000      # Queue supprimée après 30min d'inactivité
    }
)

# Gestion intelligente des TTL
class TTLManager:
    def __init__(self):
        self.ttl_strategies = {
            'critical': 3600000,    # 1 heure
            'normal': 1800000,      # 30 minutes  
            'bulk': 300000,         # 5 minutes
            'temp': 60000           # 1 minute
        }
    
    def get_ttl(self, message_type, priority=None):
        base_ttl = self.ttl_strategies.get(message_type, 1800000)
        
        if priority and priority > 8:
            return base_ttl * 2  # Plus de temps pour haute priorité
        elif priority and priority < 3:
            return base_ttl // 2  # Moins de temps pour basse priorité
        
        return base_ttl

Priority (Priorités)

// Configuration queue avec priorité
@Bean
public Queue priorityQueue() {
    return QueueBuilder
        .durable("priority.tasks")
        .withArgument("x-max-priority", 10)  // Priorité 0-10
        .build();
}

// Publishing avec priorité
@Service
public class PriorityTaskPublisher {
    
    public void publishUrgentTask(Task task) {
        rabbitTemplate.convertAndSend(
            "tasks", 
            "urgent",
            task,
            message -> {
                message.getMessageProperties().setPriority(10);  // Max priorité
                return message;
            }
        );
    }
    
    public void publishByImportance(Task task) {
        int priority = calculatePriority(task);
        
        rabbitTemplate.convertAndSend(
            "tasks",
            "standard", 
            task,
            message -> {
                message.getMessageProperties().setPriority(priority);
                message.getMessageProperties().setHeaders(Map.of(
                    "task-type", task.getType(),
                    "business-priority", task.getBusinessPriority()
                ));
                return message;
            }
        );
    }
    
    private int calculatePriority(Task task) {
        // Logique métier pour calculer priorité
        return switch (task.getType()) {
            case PAYMENT -> 9;           // Critique
            case NOTIFICATION -> 5;      // Normal
            case ANALYTICS -> 1;         // Bas
            default -> 3;
        };
    }
}

🔗 Correlation et Reply-To

Pattern Request-Reply

import uuid
import time
from threading import Event

class RPCClient:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        
        # Queue temporaire pour responses
        result = self.channel.queue_declare(queue='', exclusive=True)
        self.callback_queue = result.method.queue
        self.response = None
        self.correlation_id = None
        
        self.channel.basic_consume(
            queue=self.callback_queue,
            on_message_callback=self.on_response,
            auto_ack=True
        )
    
    def on_response(self, ch, method, props, body):
        if self.correlation_id == props.correlation_id:
            self.response = body
    
    def call(self, request_data, timeout=30):
        self.response = None
        self.correlation_id = str(uuid.uuid4())
        
        # Envoyer request avec correlation_id
        self.channel.basic_publish(
            exchange='',
            routing_key='rpc_queue',
            properties=pika.BasicProperties(
                reply_to=self.callback_queue,
                correlation_id=self.correlation_id,
                timestamp=int(time.time()),
                expiration=str(timeout * 1000)  # Timeout en ms
            ),
            body=json.dumps(request_data)
        )
        
        # Attendre response
        start_time = time.time()
        while self.response is None and (time.time() - start_time) < timeout:
            self.connection.process_data_events()
            time.sleep(0.01)
        
        if self.response is None:
            raise TimeoutError(f"RPC call timeout after {timeout}s")
        
        return json.loads(self.response)

# Usage
rpc_client = RPCClient()
result = rpc_client.call({
    'operation': 'validate_payment',
    'payment_id': 'pay_123',
    'amount': 99.99
}, timeout=10)

Tracing Distribué

@Component
public class TracingMessageProducer {
    
    private final RabbitTemplate rabbitTemplate;
    private final Tracer tracer;
    
    public void publishWithTracing(String exchange, String routingKey, Object message) {
        Span span = tracer.nextSpan()
            .name("rabbitmq.publish")
            .tag("messaging.system", "rabbitmq")
            .tag("messaging.destination", exchange)
            .tag("messaging.destination_kind", "exchange")
            .start();
        
        try (Tracer.SpanInScope ws = tracer.withSpanInScope(span)) {
            String traceId = span.context().traceId();
            String spanId = span.context().spanId();
            
            rabbitTemplate.convertAndSend(
                exchange,
                routingKey,
                message,
                msg -> {
                    msg.getMessageProperties().setHeaders(Map.of(
                        "x-trace-id", traceId,
                        "x-span-id", spanId,
                        "x-parent-span-id", span.context().parentId()
                    ));
                    return msg;
                }
            );
            
            span.tag("messaging.message_id", 
                    rabbitTemplate.getMessageId().toString());
            
        } catch (Exception e) {
            span.tag("error", e.getMessage());
            throw e;
        } finally {
            span.end();
        }
    }
}

@RabbitListener(queues = "traced.messages")
public void handleTracedMessage(
        @Payload Object message,
        @Header Map<String, Object> headers) {
    
    String traceId = (String) headers.get("x-trace-id");
    String parentSpanId = (String) headers.get("x-span-id");
    
    if (traceId != null) {
        SpanBuilder spanBuilder = tracer.nextSpan()
            .name("rabbitmq.consume")
            .tag("messaging.system", "rabbitmq");
        
        if (parentSpanId != null) {
            // Continuer la trace
            SpanContext parentContext = SpanContext.fromTraceId(
                traceId, parentSpanId, TraceFlags.getSampled()
            );
            spanBuilder.setParent(Context.root().with(Span.wrap(parentContext)));
        }
        
        Span span = spanBuilder.start();
        
        try (Tracer.SpanInScope ws = tracer.withSpanInScope(span)) {
            processMessage(message);
        } finally {
            span.end();
        }
    } else {
        processMessage(message);
    }
}

🎛️ Headers Personnalisés

Business Headers

# Headers métier riches
def publish_business_event(event_data, context):
    headers = {
        # Context business
        'tenant-id': context.tenant_id,
        'user-id': context.user_id,
        'session-id': context.session_id,
        'request-id': context.request_id,
        
        # Metadata technique  
        'service-name': 'user-service',
        'service-version': '1.2.3',
        'environment': 'production',
        'region': 'eu-west-1',
        
        # Routing hints
        'event-type': event_data['type'],
        'severity': event_data.get('severity', 'info'),
        'category': event_data.get('category', 'business'),
        
        # Processing hints
        'idempotency-key': str(uuid4()),
        'retry-policy': 'exponential',
        'max-retries': 3,
        
        # Compliance
        'data-classification': 'personal',
        'retention-policy': '7-years',
        'audit-required': True
    }
    
    channel.basic_publish(
        exchange='business.events',
        routing_key=f"user.{event_data['type']}",
        body=json.dumps(event_data),
        properties=pika.BasicProperties(
            headers=headers,
            delivery_mode=2,
            timestamp=int(time.time()),
            message_id=headers['idempotency-key']
        )
    )

Headers-Based Routing

// Exchange Headers pour routage par contenu
@Configuration
public class HeadersRoutingConfig {
    
    @Bean
    public HeadersExchange contentRouter() {
        return ExchangeBuilder
            .headersExchange("content.router")
            .durable(true)
            .build();
    }
    
    // Routage pour PDF 
    @Bean
    public Binding pdfProcessingBinding() {
        return BindingBuilder
            .bind(pdfQueue())
            .to(contentRouter())
            .whereAll(Map.of(
                "content-type", "application/pdf",
                "processing-required", true
            )).match();
    }
    
    // Routage pour images haute résolution
    @Bean 
    public Binding highResImageBinding() {
        return BindingBuilder
            .bind(imageQueue())
            .to(contentRouter())
            .whereAll(Map.of(
                "content-type", "image/jpeg",
                "resolution", "high"
            )).match();
    }
    
    // Routage conditionnel complexe
    @Bean
    public Binding urgentDocumentsBinding() {
        return BindingBuilder
            .bind(urgentQueue())
            .to(contentRouter())
            .whereAny(Map.of(
                "priority", "urgent",
                "deadline", "today"
            )).match();
    }
}

// Publisher utilisant headers routing
@Service
public class DocumentProcessor {
    
    public void processDocument(Document doc) {
        Map<String, Object> headers = Map.of(
            "content-type", doc.getMimeType(),
            "size-mb", doc.getSizeInMB(),
            "resolution", doc.isHighResolution() ? "high" : "standard",
            "priority", calculatePriority(doc),
            "tenant-id", doc.getTenantId()
        );
        
        rabbitTemplate.convertAndSend(
            "content.router",
            "", // Headers exchange ignore routing key
            doc,
            message -> {
                message.getMessageProperties().setHeaders(headers);
                return message;
            }
        );
    }
}

⏰ Timestamp et Expiration

Gestion Fine des Timeouts

class TimeoutManager:
    def __init__(self):
        self.timeout_strategies = {
            'real_time': 5000,      # 5 secondes
            'interactive': 30000,    # 30 secondes  
            'batch': 300000,        # 5 minutes
            'background': 3600000   # 1 heure
        }
    
    def publish_with_smart_timeout(self, message, processing_type):
        timeout = self.timeout_strategies.get(processing_type, 30000)
        
        # Ajuster selon la charge système
        current_load = self.get_system_load()
        if current_load > 0.8:
            timeout *= 2  # Doubler le timeout si système chargé
        
        properties = pika.BasicProperties(
            timestamp=int(time.time()),
            expiration=str(timeout),
            headers={
                'processing-type': processing_type,
                'estimated-duration': timeout,
                'max-retries': 3 if processing_type == 'critical' else 1
            }
        )
        
        channel.basic_publish(
            exchange='timed.processing',
            routing_key=processing_type,
            body=json.dumps(message),
            properties=properties
        )

Dead Letter avec TTL

// Queue avec TTL automatique vers DLX
const setupTTLQueue = async () => {
    // Queue principale avec TTL
    await channel.assertQueue('processing.main', {
        durable: true,
        arguments: {
            'x-message-ttl': 60000,           // 1 minute
            'x-dead-letter-exchange': 'retry', // DLX pour retry
            'x-dead-letter-routing-key': 'expired'
        }
    });
    
    // Queue retry avec TTL plus long
    await channel.assertQueue('processing.retry', {
        durable: true,
        arguments: {
            'x-message-ttl': 300000,          // 5 minutes
            'x-dead-letter-exchange': 'main',  // Retour vers main
            'x-dead-letter-routing-key': 'retry'
        }
    });
    
    // Queue finale pour échecs
    await channel.assertQueue('processing.failed', {
        durable: true
        // Pas de TTL - stockage permanent pour investigation
    });
};

🔄 Reply-To et Correlation ID

RPC Asynchrone Avancé

@Service
public class AsyncRPCService {
    
    private final Map<String, CompletableFuture<Object>> pendingRequests = 
        new ConcurrentHashMap<>();
    
    // Publier request avec correlation
    public CompletableFuture<PaymentResult> processPaymentAsync(PaymentRequest request) {
        String correlationId = UUID.randomUUID().toString();
        CompletableFuture<Object> future = new CompletableFuture<>();
        
        // Stocker la future pour correlation
        pendingRequests.put(correlationId, future);
        
        // Timeout automatique
        CompletableFuture.delayedExecutor(30, TimeUnit.SECONDS).execute(() -> {
            CompletableFuture<Object> timeoutFuture = pendingRequests.remove(correlationId);
            if (timeoutFuture != null && !timeoutFuture.isDone()) {
                timeoutFuture.completeExceptionally(
                    new TimeoutException("RPC timeout after 30s")
                );
            }
        });
        
        rabbitTemplate.convertAndSend(
            "payment.rpc",
            "process",
            request,
            message -> {
                message.getMessageProperties().setCorrelationId(correlationId);
                message.getMessageProperties().setReplyTo("payment.rpc.replies");
                message.getMessageProperties().setExpiration("30000");
                return message;
            }
        );
        
        return future.thenApply(result -> (PaymentResult) result);
    }
    
    // Handler des responses
    @RabbitListener(queues = "payment.rpc.replies")
    public void handleRPCResponse(
            @Payload Object response,
            @Header(AmqpHeaders.CORRELATION_ID) String correlationId) {
        
        CompletableFuture<Object> future = pendingRequests.remove(correlationId);
        if (future != null) {
            future.complete(response);
        } else {
            log.warn("Received response for unknown correlation ID: {}", correlationId);
        }
    }
}

Conversation State

# Maintenir état de conversation
class ConversationManager:
    def __init__(self):
        self.conversations = {}  # conversation_id -> state
    
    def start_conversation(self, user_id, conversation_type):
        conversation_id = str(uuid4())
        self.conversations[conversation_id] = {
            'user_id': user_id,
            'type': conversation_type,
            'state': 'started',
            'created_at': time.time(),
            'steps': []
        }
        
        return conversation_id
    
    def publish_conversation_step(self, conversation_id, step_data):
        conversation = self.conversations.get(conversation_id)
        if not conversation:
            raise ValueError(f"Unknown conversation: {conversation_id}")
        
        # Ajouter contexte de conversation
        step_data['conversation_context'] = {
            'conversation_id': conversation_id,
            'user_id': conversation['user_id'],
            'step_number': len(conversation['steps']) + 1,
            'previous_steps': conversation['steps'][-3:]  # Derniers 3 steps
        }
        
        properties = pika.BasicProperties(
            correlation_id=conversation_id,
            headers={
                'conversation-id': conversation_id,
                'conversation-type': conversation['type'],
                'step-number': step_data['conversation_context']['step_number'],
                'user-id': conversation['user_id']
            },
            timestamp=int(time.time())
        )
        
        channel.basic_publish(
            exchange='conversations',
            routing_key=f"conversation.{conversation['type']}.step",
            body=json.dumps(step_data),
            properties=properties
        )
        
        # Mettre à jour état
        conversation['steps'].append(step_data)
        conversation['last_activity'] = time.time()

🎯 Patterns Avancés

Message Deduplication

// Déduplication basée sur message ID
@Component
public class DeduplicationService {
    
    private final RedisTemplate<String, String> redis;
    private static final Duration DEDUP_TTL = Duration.ofHours(24);
    
    @RabbitListener(queues = "dedup.payments")
    public void handlePayment(
            @Payload PaymentMessage payment,
            @Header(AmqpHeaders.MESSAGE_ID) String messageId) {
        
        String dedupKey = "processed:" + messageId;
        
        // Check if already processed
        Boolean isNew = redis.opsForValue().setIfAbsent(
            dedupKey, 
            "processed", 
            DEDUP_TTL
        );
        
        if (Boolean.TRUE.equals(isNew)) {
            // First time processing this message
            try {
                paymentService.processPayment(payment);
                
                // Mark as successfully processed
                redis.opsForValue().set(
                    "success:" + messageId,
                    payment.toString(),
                    DEDUP_TTL
                );
                
            } catch (Exception e) {
                // Remove dedup key on failure to allow retry
                redis.delete(dedupKey);
                throw e;
            }
        } else {
            log.info("Duplicate message detected, skipping: {}", messageId);
            
            // Optional: vérifier si traitement était successful
            String result = redis.opsForValue().get("success:" + messageId);
            if (result == null) {
                log.warn("Previous processing may have failed, consider retry");
            }
        }
    }
}

Batch Processing avec Properties

# Accumulation et traitement par batch
class BatchProcessor:
    def __init__(self, batch_size=100, timeout_seconds=30):
        self.batch_size = batch_size
        self.timeout_seconds = timeout_seconds
        self.current_batch = []
        self.batch_start_time = None
    
    def process_message(self, ch, method, properties, body):
        message = json.loads(body)
        
        # Ajouter au batch
        self.current_batch.append({
            'message': message,
            'properties': {
                'message_id': properties.message_id,
                'timestamp': properties.timestamp,
                'headers': properties.headers or {}
            },
            'delivery_tag': method.delivery_tag
        })
        
        if self.batch_start_time is None:
            self.batch_start_time = time.time()
        
        # Décider si traiter le batch
        should_process = (
            len(self.current_batch) >= self.batch_size or
            (time.time() - self.batch_start_time) > self.timeout_seconds or
            properties.headers.get('force-batch-process') == 'true'
        )
        
        if should_process:
            self.process_current_batch(ch)
    
    def process_current_batch(self, channel):
        if not self.current_batch:
            return
        
        try:
            # Traitement batch optimisé
            results = self.business_service.process_batch([
                item['message'] for item in self.current_batch
            ])
            
            # Ack tous les messages du batch
            for item in self.current_batch:
                channel.basic_ack(
                    delivery_tag=item['delivery_tag'],
                    multiple=False
                )
            
            # Publier résultats
            self.publish_batch_results(results)
            
        except Exception as e:
            # Nack tout le batch pour retry
            for item in self.current_batch:
                channel.basic_nack(
                    delivery_tag=item['delivery_tag'],
                    multiple=False,
                    requeue=True
                )
            raise
        
        finally:
            # Reset batch
            self.current_batch = []
            self.batch_start_time = None

📊 Monitoring des Properties

Analyse des Patterns

// Analyser les patterns d'utilisation
class MessagePropertiesAnalyzer {
    constructor() {
        this.metrics = new Map();
    }
    
    analyzeMessage(properties, body) {
        const analysis = {
            size: Buffer.byteLength(body),
            contentType: properties.contentType,
            priority: properties.priority || 0,
            persistent: properties.deliveryMode === 2,
            hasHeaders: !!(properties.headers && Object.keys(properties.headers).length),
            age: Date.now() - (properties.timestamp * 1000)
        };
        
        // Collecter métriques
        this.updateMetrics(analysis);
        
        // Détecter anomalies
        this.detectAnomalies(analysis, properties);
        
        return analysis;
    }
    
    detectAnomalies(analysis, properties) {
        // Messages très anciens
        if (analysis.age > 300000) { // 5 minutes
            console.warn('Old message detected:', {
                messageId: properties.messageId,
                age: analysis.age,
                headers: properties.headers
            });
        }
        
        // Messages très gros
        if (analysis.size > 1024 * 1024) { // 1MB
            console.warn('Large message detected:', {
                messageId: properties.messageId,
                size: analysis.size,
                contentType: analysis.contentType
            });
        }
        
        // Priorité incohérente
        if (properties.headers?.['business-priority'] === 'low' && 
            analysis.priority > 5) {
            console.warn('Priority mismatch:', properties);
        }
    }
}

✅ Bonnes Pratiques

1. Standardisation des Properties

# Schema standardisé pour l'organisation
MESSAGE_SCHEMA = {
    'required_headers': [
        'service-name',
        'service-version', 
        'tenant-id',
        'trace-id'
    ],
    'optional_headers': [
        'user-id',
        'session-id',
        'business-category',
        'data-classification'
    ],
    'reserved_headers': [
        'x-retry-count',
        'x-original-queue',
        'x-death'  # RabbitMQ internal
    ]
}

def validate_message_properties(properties):
    headers = properties.headers or {}
    
    # Vérifier headers requis
    for required in MESSAGE_SCHEMA['required_headers']:
        if required not in headers:
            raise ValueError(f"Missing required header: {required}")
    
    # Vérifier headers réservés
    for reserved in MESSAGE_SCHEMA['reserved_headers']:
        if reserved in headers and not reserved.startswith('x-'):
            raise ValueError(f"Cannot use reserved header: {reserved}")
    
    return True

2. Performance Guidelines

# Recommandations de performance
performance_guidelines:
  message_size:
    recommended_max: "64KB"
    absolute_max: "128MB" 
    consideration: "Large messages impact memory and performance"
  
  headers:
    recommended_max: 10
    size_limit: "4KB total"
    consideration: "Too many headers slow down routing"
    
  priority:
    use_sparingly: true
    max_levels: 5
    consideration: "Priority queues have overhead"
    
  ttl:
    always_set: true
    default: "30 minutes"
    consideration: "Prevents infinite accumulation"

Maîtriser les propriétés des messages est essentiel pour exploiter pleinement la flexibilité et la puissance de RabbitMQ dans vos architectures.

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours