📦 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.