Pular para o conteúdo

Relay AMQP

Guia de uso do canal de entrega contínua. A referência completa do canal está na navegação lateral, em Canais assíncronos.

exchange: weather.readings (topic, durable)
routing key: station.{stationId}.reading
┌───────────────┴───────────────┐
▼ ▼
relay.acme (sua fila) relay.outro-parceiro
seu consumidor ──rejeita──▶ relay.acme.dlq

A fila é sua, é durável e é criada por nós. Suas credenciais dão consumo apenas: você não declara filas, não altera bindings e não publica. A restrição protege os outros integradores de uma configuração equivocada em qualquer cliente.

Para vincular ou desvincular estações da sua fila, fale conosco: o binding é feito do nosso lado.

amqps://relay.radar-sandbox.kitelife.com.br:5671/integradores

TLS obrigatório. A 5672 em texto claro não é exposta. Usuário, senha e o nome da sua fila vêm no provisionamento.

import amqp from 'amqplib'
const conn = await amqp.connect(process.env.RADAR_AMQP_URL)
const ch = await conn.createChannel()
await ch.prefetch(100)
await ch.consume('relay.acme', async (msg) => {
const leitura = JSON.parse(msg.content.toString())
await gravar(leitura) // idempotente, ver abaixo
ch.ack(msg) // só depois de gravar
}, { noAck: false })

A entrega é ao menos uma vez, e o message_id é determinístico: a mesma leitura da mesma estação sempre produz o mesmo identificador, derivado de stationId e observedAt. Deduplique por ele, ou, melhor ainda, torne a escrita idempotente no banco:

CREATE UNIQUE INDEX ON readings (station_id, observed_at);
-- e então
INSERT INTO readings (…) VALUES (…) ON CONFLICT DO NOTHING;

Assim a duplicata custa uma linha ignorada em vez de um ponto duplicado na série. Deduplicar em memória, com um conjunto de ids vistos, funciona até o seu processo reiniciar.

Dentro de uma estação, as mensagens saem em ordem cronológica de observedAt. Entre estações não há ordenação: cada uma publica no seu ritmo.

Com a cadência das estações, um prefetch entre 50 e 200 costuma ser o ponto de equilíbrio para um consumidor único.

A fila acumula. Trinta estações a cada oito segundos são cerca de 300 mil mensagens por dia, e a fila tem limite de tamanho e tempo de vida: passado o limite, as mensagens mais antigas são descartadas para dar lugar às novas.

Isso é preferível a estourar o broker, mas significa que uma parada longa deixa buraco. Para preencher, use a REST com from na última leitura que você gravou.

Se a sua janela de manutenção é previsível e longa, avise: dá para aumentar o limite da sua fila temporariamente.

Mensagem rejeitada sem reenfileiramento vai para {sua-fila}.dlq, que retém por 7 dias.

try {
await gravar(leitura)
ch.ack(msg)
} catch (erro) {
const reenfileirar = !(erro instanceof SchemaError)
ch.nack(msg, false, reenfileirar)
}

Além das leituras, a fila recebe station.status quando uma estação vinculada muda de estado. É assim que você distingue “a estação caiu” de “meu consumidor travou”. Sem esse evento, os dois casos são idênticos vistos de dentro do consumidor: mensagens que param de chegar.