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: direct publica 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: outbox grava a mensagem de forma durável dentro da mesma transação de banco da rota (exige transaction: 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:
        - ...
  • source referencia um broker declarado em spec.brokers.
  • queue (RabbitMQ) ou topic/consumerGroup (Kafka) identificam de onde consumir.
  • retry/deadLetter controlam o que acontece quando o pipeline falha ao processar uma mensagem.
  • idempotency evita reprocessar a mesma mensagem duas vezes: key é uma expressão que identifica a mensagem de forma única, store referencia onde essa chave é registrada, ttl por 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-stack no repositório do Pipevine.