Mensageria e eventos
Declarando brokers e eventos
spec:
brokers:
- { id: rabbit-main, provider: rabbitmq, connectionString: "${config:Rabbit}" }
- { id: kafka-main, provider: kafka, bootstrapServers: "${config:Kafka}" }
events:
- name: OrderCreated
schema: schemas/order-created.json
broker: kafka-main
topic: orders.v1
partitionKey: "payload.customerId"
description: Emitido quando um pedido é aceito
provider: rabbitmq e provider: kafka são as implementações de referência. Um broker novo pode
ser adicionado por um pacote de terceiro — veja
IBrokerProvider. Cada evento declarado alimenta a
geração automática de AsyncAPI, sem nenhuma anotação manual.
Publicando: o step publish
- type: publish
event: OrderCreated
payload: "vars.order"
mode: outbox # direct | outbox
event: referencia um item de spec.events; payload: é uma expressão JSONata que monta o corpo
da mensagem.
mode: directpublica no broker de forma síncrona, dentro do próprio step. Simples, mas se a transação da rota for revertida depois da publicação, o evento já foi enviado — não há como desfazer.mode: outboxgrava a mensagem de forma durável dentro da mesma transação de banco da rota (exigetransaction: required), e um despachante em segundo plano entrega a mensagem depois, com garantia de pelo menos uma entrega (at-least-once). Isso é o que garante que um evento nunca é publicado sem que a escrita correspondente exista de fato — e vice-versa: se a transação reverte, o evento nunca sai.
Prefira mode: outbox sempre que a publicação estiver associada a uma escrita no banco que precisa
acontecer atomicamente com ela — é o caso mais comum (publicar OrderCreated junto da persistência
do pedido, por exemplo).
Consumindo: spec.type: consumer
Uma aplicação do tipo consumer compartilha exatamente o mesmo modelo de pipeline de uma api — os
steps disponíveis dentro do pipeline de um consumidor são os mesmos que você já viu nos guias
anteriores.
spec:
type: consumer
consumers:
- id: onOrderPaid
source: rabbit-main # referência a spec.brokers
queue: orders.paid # RabbitMQ
# topic: orders.paid # Kafka
# consumerGroup: orders-service
concurrency: 4
prefetch: 20
retry:
attempts: 5
backoff: { initial: 1s, max: 30s, jitter: true }
deadLetter:
queue: orders.paid.dlq
idempotency:
key: "input.headers.'ce-id'"
store: postgres-main
ttl: 24h
pipeline:
- ...
sourcereferencia um broker declarado emspec.brokers.queue(RabbitMQ) outopic/consumerGroup(Kafka) identificam de onde consumir.retry/deadLettercontrolam o que acontece quando o pipeline falha ao processar uma mensagem.idempotencyevita reprocessar a mesma mensagem duas vezes:keyé uma expressão que identifica a mensagem de forma única,storereferencia onde essa chave é registrada,ttlpor quanto tempo a chave é lembrada.
A chave de idempotência precisa ser única por consumidor, não só por entidade de negócio. Se
dois consumidores diferentes processam eventos relacionados ao mesmo pedido e ambos usam só o id do
pedido como chave, o segundo consumidor a ver aquele id trata como duplicata e pula o processamento
silenciosamente — mesmo sendo uma mensagem legítima de um tópico diferente. Prefixar a chave com o
nome do próprio consumidor (ou, melhor ainda, usar o id da própria mensagem — um ce-id de
CloudEvents, por exemplo, que já é único por mensagem) evita esse problema.
Para onde ir daqui
- Integrações outbound
- Um exemplo real com dois brokers e outbox:
samples/06-full-stackno repositório do Pipevine.