| name | async-messaging-patterns |
| description | Padrões de mensageria assíncrona e event-driven: SQS, SNS, Kafka, EventBridge, idempotência, DLQ, ordering, retry. Use quando implementar consumers, producers ou arquitetura event-driven. |
| argument-hint | [contexto adicional] |
Async Messaging — Patterns & Idioms
Padroes para mensageria assincrona e arquiteturas event-driven em producao.
Escolha de tecnologia
| Criterio | SQS | SNS+SQS | Kafka | EventBridge |
|---|
| Point-to-point | Sim | — | — | — |
| Fan-out (1→N) | — | Sim | Sim (consumer groups) | Sim (rules) |
| Ordering | FIFO (por group) | FIFO | Sim (por partition) | Nao garantido |
| Replay | Nao | Nao | Sim (retention) | Archive+Replay |
| Throughput | Alto | Alto | Muito alto | Moderado |
| Latencia | ~ms | ~ms | ~ms | ~50-100ms |
| Serverless native | Sim (Lambda trigger) | Sim | Nao (MSK) | Sim |
| Custo em idle | Zero | Zero | Fixo (brokers) | Zero |
Estrutura de evento
{
"eventId": "evt-abc123",
"eventType": "order.created",
"eventVersion": "1.0",
"source": "order-service",
"timestamp": "2026-04-18T10:30:00Z",
"correlationId": "req-xyz789",
"data": {
"orderId": "ord-456",
"customerId": "cust-789",
"amount": 299.90,
"status": "pending"
},
"metadata": {
"traceId": "trace-abc",
"environment": "production"
}
}
Regras de design de eventos
- eventId: UUID unico por evento — essencial para idempotencia
- eventType:
<domain>.<action> em past tense (order.created, nao create.order)
- eventVersion: versionamento explicito do schema
- correlationId: rastrear fluxo end-to-end
- data: payload do evento — nao incluir dados sensiveis
- Eventos sao fatos imutaveis — nao comandos
Idempotencia (obrigatorio)
Todo consumer DEVE ser idempotente. Mensagens podem ser entregues mais de uma vez.
@Service
public class IdempotentProcessor {
private final DynamoDbTable<ProcessedEvent> table;
public boolean tryProcess(String eventId, Runnable action) {
try {
table.putItem(PutItemEnhancedRequest.builder(ProcessedEvent.class)
.item(new ProcessedEvent(eventId, Instant.now()))
.conditionExpression(Expression.builder()
.expression("attribute_not_exists(eventId)")
.build())
.build());
action.run();
return true;
} catch (ConditionalCheckFailedException e) {
log.info("Event already processed: {}", eventId);
return false;
}
}
}
from aws_lambda_powertools.utilities.idempotency import (
idempotent, DynamoDBPersistenceLayer, IdempotencyConfig,
)
persistence = DynamoDBPersistenceLayer(table_name="idempotency")
config = IdempotencyConfig(event_key_jmespath="detail.eventId", expires_after_seconds=3600)
@idempotent(config=config, persistence_store=persistence)
def process_order(event):
order_id = event["detail"]["data"]["orderId"]
order_service.process(order_id)
func (p *Processor) Process(ctx context.Context, event Event) error {
processed, err := p.idempotencyStore.Exists(ctx, event.EventID)
if err != nil { return fmt.Errorf("idempotency check: %w", err) }
if processed { return nil }
if err := p.service.Handle(ctx, event); err != nil { return err }
return p.idempotencyStore.Mark(ctx, event.EventID)
}
SQS + Lambda
public class OrderSqsHandler implements RequestHandler<SQSEvent, SQSBatchResponse> {
@Override
public SQSBatchResponse handleRequest(SQSEvent event, Context context) {
var failures = new ArrayList<SQSBatchResponse.BatchItemFailure>();
for (var record : event.getRecords()) {
try {
var orderEvent = objectMapper.readValue(record.getBody(), OrderEvent.class);
idempotentProcessor.tryProcess(orderEvent.eventId(),
() -> orderService.process(orderEvent));
} catch (Exception e) {
log.error("Failed to process message: {}", record.getMessageId(), e);
failures.add(new SQSBatchResponse.BatchItemFailure(record.getMessageId()));
}
}
return new SQSBatchResponse(failures);
}
}
OrderFunction:
Type: AWS::Serverless::Function
Properties:
Events:
SQSEvent:
Type: SQS
Properties:
Queue: !GetAtt OrderQueue.Arn
BatchSize: 10
MaximumBatchingWindowInSeconds: 5
FunctionResponseTypes:
- ReportBatchItemFailures
Kafka (consumer)
@Component
public class OrderKafkaConsumer {
@KafkaListener(
topics = "orders",
groupId = "order-processor",
containerFactory = "kafkaListenerContainerFactory",
)
public void consume(
@Payload String payload,
@Header(KafkaHeaders.RECEIVED_KEY) String key,
@Header("eventId") String eventId,
Acknowledgment ack
) {
try {
var event = objectMapper.readValue(payload, OrderEvent.class);
if (idempotentProcessor.tryProcess(eventId, () -> orderService.process(event))) {
ack.acknowledge();
} else {
ack.acknowledge();
}
} catch (Exception e) {
log.error("Failed to process event: {}", eventId, e);
throw e;
}
}
}
func (c *Consumer) Run(ctx context.Context) error {
for {
msg, err := c.reader.ReadMessage(ctx)
if err != nil {
if errors.Is(err, context.Canceled) { return nil }
return fmt.Errorf("read message: %w", err)
}
var event OrderEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
slog.Error("invalid message", "error", err, "offset", msg.Offset)
continue
}
if err := c.processor.Process(ctx, event); err != nil {
slog.Error("process failed", "error", err, "eventId", event.EventID)
}
}
}
EventBridge
public class OrderEventPublisher {
private final EventBridgeClient eventBridge;
public void publish(OrderEvent event) {
eventBridge.putEvents(PutEventsRequest.builder()
.entries(PutEventsRequestEntry.builder()
.source("order-service")
.detailType("order.created")
.detail(objectMapper.writeValueAsString(event))
.eventBusName("orders")
.build())
.build());
}
}
def handler(event, context):
detail = event["detail"]
event_id = detail["eventId"]
order_id = detail["data"]["orderId"]
logger.info("processing", event_id=event_id, order_id=order_id)
order_service.process(detail)
Retry com backoff e jitter
OrderQueue:
Type: AWS::SQS::Queue
Properties:
VisibilityTimeout: 300
RedrivePolicy:
deadLetterTargetArn: !GetAtt OrderDLQ.Arn
maxReceiveCount: 3
OrderDLQ:
Type: AWS::SQS::Queue
Properties:
MessageRetentionPeriod: 1209600
@Retryable(
maxAttempts = 3,
backoff = @Backoff(delay = 1000, multiplier = 2, maxDelay = 10000, random = true)
)
public void publishEvent(OrderEvent event) {
eventBridge.putEvents();
}
DLQ e poison message handling
def dlq_handler(event, context):
for record in event["Records"]:
body = json.loads(record["body"])
logger.warning("DLQ message",
event_id=body.get("eventId"),
error=record.get("attributes", {}).get("DeadLetterQueueSourceArn"),
receive_count=record["attributes"]["ApproximateReceiveCount"],
)
metrics.increment("dlq.messages.received", tags={"source": body.get("source")})
Schema evolution
| Estrategia | Descricao | Quando |
|---|
| Additive | Adicionar campos novos (nullable/default) | Sempre preferido |
| Versioned | Novo eventType com versao (order.created.v2) | Breaking change necessario |
| Schema Registry | Avro/Protobuf com compatibilidade BACKWARD | Kafka com contratos fortes |
currency = event.get("currency", "BRL")
Observabilidade
var traceId = event.metadata().traceId();
Span span = tracer.spanBuilder("process-order")
.setParent(Context.current().with(extractTraceContext(traceId)))
.startSpan();
Checklist