Skip to main content
Um evento percorre o Streaming Hub em três etapas. O hub o ingere do stream interno, o compara com as suas subscriptions e o despacha para cada destino correspondente. Esta página percorre cada etapa e as garantias de entrega que saem dela.

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:
O 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 um origin fixo na subscription quando a mesma chave puder vir de mais de uma aplicação.
Veja a visão geral de streaming de eventos para o contrato de fio completo.

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:
A curva cobre cerca de 28 horas ao longo dos seus oito degraus. Cada subscription pode sobrescrevê-la com um agendamento próprio (até 12 degraus, cada um limitado a 10 horas). Uma sobrescrita malformada volta para a curva padrão em vez de fazer a entrega falhar.

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.
Os dois portões devem valer, e uma única entrega bem-sucedida limpa o período, então uma oscilação breve nunca acumula rumo ao veredito. Quando a desabilitação automática dispara, ela vira a flag 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.