Ingerir eventos do stream
O Streaming Hub consome o stream da plataforma como mensagens CloudEvents 1.0 em modo de conteúdo binário. Os atributos de contexto do CloudEvents viajam nos headers do registro Kafka (cada um com o prefixo
ce-). O corpo do evento é o valor do registro. O hub lê roteamento e identidade dos headers sem desserializar o payload.
O hub assina apenas os tópicos de fato de aplicação v3:
ce-source é um segmento único, em minúsculas e sem pontos. Ele usa letras, dígitos, sublinhados ou hifens. O tópico não carrega sufixo de recurso, de evento nem de major do schema. O hub despacha pelos headers do CloudEvents. O conjunto seguido exclui tópicos de comando (.commands) e tópicos de dead-letter (.dlq). Os tópicos de fato de aplicação do Midaz, do Tracer, do Matcher, do Lender e do Consignado agora atendem a esse contrato de conjunto seguido.
Um segundo consumidor lê os tópicos .dlq que os produtores gravam. Esse consumidor serve apenas à observabilidade. A gramática do conjunto seguido exclui um .dlq final, então um registro de dead-letter nunca vira uma entrega para uma das suas subscriptions. Os operadores leem esses registros por GET /admin/dlq. Veja forense de DLQ.
Cada registro que o hub aceita deve carregar este conjunto de contexto:
Para cada registro o hub devolve um de três vereditos:
O tenant admitido na ingestão vem apenas do atributo
ce-tenantid validado que o parser extraiu. A ingestão não tem portão de descarte por lista de tenants. Os eventos guardados apenas correspondem a subscriptions com o mesmo id de tenant.
Três strings parecidas são coisas distintas. Nunca as confunda:
ce-typeé o tipo CloudEvents qualificado pela origem:studio.lerian.<source>.<resource>.<event>(por exemplo,studio.lerian.lender.loan_application.approved).- O tópico Kafka é um tópico de fato v3 por aplicação:
lerian.streaming.<source>(por exemplo,lerian.streaming.lender). Recurso, evento e major do schema não aparecem no tópico. X-Lerian-Event-Typeé a chave<resource>.<event>pura em um webhook entregue (loan_application.approved). Ela omite o produtor de propósito. Use umoriginfixo na subscription quando a mesma chave puder vir de mais de uma aplicação.
Deduplicação de consumo único
O stream interno é pelo menos uma vez, então o mesmo registro pode chegar mais de uma vez. No caminho para a caixa de entrada de eventos, o hub deduplica pelo
ce-id. Um id duplicado não grava nada e conta como descarte por dedup. O hub grava registros por partição Kafka em uma transação. Ele confirma a posição no stream de uma partição apenas depois que a gravação daquela partição é confirmada. Uma falha em uma partição, portanto, nunca perde nem confirma duas vezes eventos de outra.
Correspondência de subscriptions
Depois que o hub guarda um evento, ele resolve quais das suas subscriptions devem recebê-lo. A correspondência de subscriptions é interna e livre de catálogo: ela avalia cada evento diretamente por
(tenant, origin, event type, schema major). A dimensão origin é opcional na subscription. Sem ela, a mesma chave de evento vinda de qualquer aplicação produtora pode corresponder. Um evento que nenhuma subscription quer produz zero entregas. Esse é um resultado correto, não um erro.
A correspondência admite uma subscription apenas quando ela está enabled e em verification_state = active. Essas são duas condições independentes, e a correspondência exige as duas. Veja o modelo de subscription.
Despachar para os seus destinos
Para cada correspondência, o hub cria um job de entrega e o passa a um pool de workers (oito workers por padrão). Os workers reivindicam jobs vencidos de forma atômica com um lease curto e intercalam entre tenants, então o acúmulo de um único tenant nunca deixa os outros sem vez. Cada worker carrega o payload guardado, descriptografa em memória o signing secret ou a credencial do destino, e faz exatamente uma tentativa de entrega ao sink. A entrega é pelo menos uma vez: o hub pode entregar o mesmo evento mais de uma vez (por novas tentativas ou reentregas depois do reinício de um worker). Cada entrega carrega um
X-Lerian-Event-Id estável (o ce-id) para você deduplicar, e um X-Lerian-Delivery-Id por tentativa que muda a cada nova tentativa. Trate um X-Lerian-Event-Id repetido como duplicata e o confirme sem processar de novo. Veja Consumir eventos para o lado do consumidor desse contrato.
Novas tentativas e backoff
Quando uma tentativa falha, o hub agenda a próxima em uma curva de backoff fixa com jitter completo:
Dead-letter
O hub manda para dead-letter um evento que esgota suas novas tentativas. O hub para de tentá-lo e registra o resultado terminal. O hub recupera um job que, em vez disso, derruba um worker no meio da tentativa e nunca registra um resultado. O limite de recuperação é cinco por padrão. O hub então o manda para dead-letter como venenoso, assim um único job tóxico nunca ocupa um worker para sempre.
Circuit breaker
Cada par
(tenant, destination) tem um circuit breaker. Quando um destino falha repetidamente, o breaker dele abre. Os jobs seguintes para aquele destino são dispensados (reagendados sem tentativa). Um único endpoint quebrado não queima capacidade de worker nem martela um alvo em dificuldade. O breaker se recupera sozinho assim que o destino volta a ter sucesso.
Desabilitar automaticamente um destino quebrado
O hub desabilita automaticamente um destino que fica quebrado por muito tempo. A desabilitação automática dispara apenas quando o período de falha aberto é os dois:
- contínuo: dura pelo menos a janela de falha (120 horas por padrão), e
- espalhado: pelo menos 12 horas separam a primeira e a última falha do período.
enabled da subscription para false. Ela nunca toca no verification_state. A entrega para (a correspondência exige enabled), e o hub registra o motivo.
Você recupera uma subscription desabilitada automaticamente com POST /v1/subscriptions/:id/verify, que sonda o destino de novo e, em caso de sucesso, a reabilita no lugar. Um kill switch (STREAMING_HUB_AUTODISABLE_ENABLED) governa a desabilitação automática, então um operador pode entregá-la desligada. Veja recuperar uma subscription desabilitada automaticamente.
Garantias de ordenação
O stream interno é particionado por tenant. Os eventos de um tenant normalmente viajam em uma única partição, e o hub preserva a ordem primeiro a entrar, primeiro a sair de ponta a ponta. O hub atribui e deduplica pelos headers do CloudEvents, não pela chave do registro Kafka. Por isso, um produtor que espalha (salga) um tenant quente por várias partições muda apenas a colocação física. O hub continua atribuindo e deduplicando corretamente. A ordenação é a única troca. Os eventos de um tenant salgado se espalham por várias partições sem garantia de ordem entre partições, então esse tenant abre mão do FIFO estrito. A sequência de chegada que o hub atribui reflete a ordem em que os eventos chegaram ao hub, não a ordem em que foram produzidos. Um tenant que não é salgado mantém o FIFO de partição única do começo ao fim.
Próximos passos
Gerenciar subscriptions
O modelo de subscription, os fluxos de onboarding e a rotação de segredos.
Consumir eventos
Verifique assinaturas, deduplique e puxe eventos.

