Documentation / Kafka Streams

07

Stateless Processing

Routing, keys, rebalances e poison records na topologia sem state store.

LAB-VERIFIEDLABKAFKA 4.3.1

Stateless Processing

O lab kafka-streams-core validou uma topologia sem banco, state store, aggregate ou join:

payment.created (Avro, key=userId)
  → validate-payment
  → VALID   → normalize-payment   → payment.validated
  → INVALID → rejection event     → payment.rejected

Routing e invariantes

Os eventos válidos preservaram eventId, userId, paymentId, amount, currency e timestamp. A regra multi-violação é determinística: MISSING_EVENT_ID, MISSING_PAYMENT_ID, MISSING_USER_ID, INVALID_AMOUNT, MISSING_CURRENCY, INVALID_TIMESTAMP; é emitida a primeira violação nessa ordem.

Key preservation e group

Não houve selectKey, groupBy ou repartition. A key original chegou intacta aos dois sinks. Um evento sem userId conservou a key vazia, porque não foi inventada uma alternativa semântica. O application.id=kafkapay-payment-validation-v1 identifica o group e coordena offsets.

Duas instâncias e failover

Duas instâncias com o mesmo application ID dividiram as seis partitions. Uma foi terminada com SIGKILL; a sobrevivente recebeu as partitions depois do rebalance e processou a carga restante. O group terminou com CURRENT-OFFSET = LOG-END-OFFSET e LAG=0.

Backlog e poison

Records publicados sem instância ativa permaneceram no input topic e foram processados no restart. Bytes sem wire format Avro foram tratados por LogAndContinueExceptionHandler: o erro foi registado, o record foi deliberadamente perdido e os seguintes continuaram. O DLQ nativo (payment.validation.dlq) foi também comprovado para desserialização; isto é diferente de payment.rejected, que contém um evento Avro válido com falha de negócio.

ALO

No ensaio isolado, o output foi encaminhado antes de kill -9, mas o offset ainda não estava commitado. O restart reprocessou o mesmo input e produziu outro output. ALO é uma política de progresso, não uma deduplicação de aplicação; se o efeito exigir unicidade, é necessário eventId e uma fronteira idempotente.