Skip to main content
Todo job de extração termina em um único evento terminal. O Worker publica job.completed quando o resultado chega ao armazenamento, e job.failed quando a execução para. Assine esses eventos em vez de consultar GET /v1/fetcher/{id} periodicamente. Esta página cobre o lado do consumidor: o que chega, o que significa e o que fazer com isso.

Onde os eventos chegam


O Worker publica em um exchange topic durável. RABBITMQ_JOB_EVENTS_EXCHANGE o nomeia, e o valor distribuído é fetcher.job.events. Vincule a sua própria fila às routing keys que lhe interessam. A definição de infraestrutura local vincula duas filas como exemplo: fetcher.job.completed.queue e fetcher.job.failed.queue.
Os eventos de job são um contrato do produto, não um recurso opcional. O Worker recusa iniciar sem o streaming habilitado e sem um nome de exchange. Veja Implantação para o lado do operador.

job.completed


Este evento significa uma coisa: a extração rodou até o fim, e o resultado criptografado já está no armazenamento de objetos. O Worker o publica depois de dois passos anteriores. Primeiro ele grava o objeto de resultado, depois registra o status terminal no job. Só então emite o evento.
metadata carrega os metadados que você enviou no job, então source e qualquer identificador de correlação voltam para você sem alteração. sizeBytes e rowCount descrevem o resultado em texto claro, antes da criptografia. path é a chave do objeto de resultado armazenado.

job.failed


Este evento significa que a execução parou e o Fetcher não armazenou resultado algum. A extração falha rápido, então a primeira fonte de dados que falha encerra o job inteiro. Um resultado parcial nunca chega ao armazenamento. O evento dispara para qualquer falha ao longo do caminho: um nome de conexão que não resolve, um erro da fonte de dados, um esquema incompatível ou uma gravação no armazenamento que não se completou.
Um evento de falha não carrega bloco result nem completedAt. metadata.error.message passa antes por uma etapa de ocultação. O Fetcher troca quatro formatos de vazamento por [redacted]: URIs de conexão, o operando de endereço de um erro de rede do Go, o operando Addr: de um erro do driver do MongoDB e um literal IPv4. O restante do texto permanece, então a mensagem continua acionável. Encaminhe-a para os operadores.

O envelope CloudEvents


Toda mensagem viaja em modo binário do CloudEvents, versão 1.0. Os atributos de contexto vão como cabeçalhos AMQP. Uma implantação single-tenant também carrega um valor de tenant. Ela emite o literal single-tenant, de modo que um mesmo consumidor atende às duas formas de implantação com o mesmo código.

Configuração do source de CloudEvents


STREAMING_CLOUDEVENTS_SOURCE não tem valor padrão. O Worker o exige sempre que o streaming está ligado, e para na inicialização quando o valor está vazio. O exemplo distribuído usa //lerian.fetcher/worker. O Fetcher copia o valor para ce-source sem alteração. Dê a cada implantação do Worker o seu próprio valor de source quando vários produtores compartilharem um broker, e roteie por esse cabeçalho.

Contrato de entrega


A entrega é at-least-once. Deduplique por ce-id.
  • O Worker grava o evento em um outbox durável antes de publicar. Uma queda do broker atrasa o evento, não o perde.
  • Um reparador varre eventos terminais que nunca foram publicados e os reemite a cada 30 segundos.
  • O ce-id permanece idêntico em toda reemissão do mesmo job e status. Nada mais é estável o bastante para servir de chave.
  • A ordem não é garantida. Dois jobs podem terminar em uma ordem e chegar em outra.
  • O Worker registra o status terminal do job antes de emitir. GET /v1/fetcher/{id} continua sendo a autoridade sobre o estado do job.
Trate um ce-id repetido como uma duplicata e confirme-o sem reprocessar. Um consumidor que usa como chave o identificador de mensagem do broker vai processar o mesmo job duas vezes.

Verificar o que você recebe


Duas assinaturas independentes protegem o caminho entre o Worker e você. As duas usam HMAC-SHA256 com a chave HMAC externa. HKDF-SHA256 deriva essa chave da chave-mestra APP_ENC_KEY. O repositório do Fetcher traz uma pequena ferramenta que imprime essa chave, de modo que um consumidor verifica as assinaturas sem nunca guardar a chave-mestra. A mensagem. O Worker assina cada mensagem publicada e carimba três cabeçalhos nela: x-message-signature, t para o timestamp da assinatura, e signature-version. O payload assinado amarra o timestamp, a versão da assinatura, o tenant, o identificador do job, o exchange e a routing key, seguidos do corpo da mensagem. Um replay do mesmo corpo sob outro tenant ou outra rota falha na verificação. O resultado. result.hmac é o HMAC-SHA256 com chave sobre o JSON do resultado em texto claro, calculado antes da criptografia. O bloco integrity declara o mesmo valor junto com o algoritmo. Verifique-o depois de descriptografar e antes de confiar nas linhas. O bloco protection descreve os bytes armazenados: encrypted é verdadeiro, o adaptador de armazenamento aplicou a criptografia, e o modo é adapter-managed. Ele descreve apenas o resultado, nunca as credenciais da fonte de dados.

Próximos passos


Jobs de extração

O que um job pede, e os estados pelos quais ele passa.

API REST do Fetcher

Criar um job, ler um job e gerenciar conexões.

Configuração

As variáveis de streaming, exchange e criptografia por trás desses eventos.

Arquitetura

O Manager, o Worker e o Engine que os dois executam.