Created
September 18, 2026 13:18
-
-
Save rponte/27dd0b4bb03b4a13cc910f5aec555b68 to your computer and use it in GitHub Desktop.
A naive idempotent Kafka consumer written in Java and Spring Boot
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package com.example.ledger; | |
| import org.apache.kafka.clients.consumer.ConsumerRecord; | |
| import org.slf4j.Logger; | |
| import org.slf4j.LoggerFactory; | |
| import org.springframework.kafka.annotation.KafkaListener; | |
| import org.springframework.kafka.support.Acknowledgment; | |
| import org.springframework.stereotype.Component; | |
| import software.amazon.awssdk.services.dynamodb.DynamoDbClient; | |
| import software.amazon.awssdk.services.dynamodb.model.AttributeValue; | |
| import software.amazon.awssdk.services.dynamodb.model.ConditionalCheckFailedException; | |
| import software.amazon.awssdk.services.dynamodb.model.GetItemRequest; | |
| import software.amazon.awssdk.services.dynamodb.model.GetItemResponse; | |
| import software.amazon.awssdk.services.dynamodb.model.PutItemRequest; | |
| import java.util.HashMap; | |
| import java.util.Map; | |
| /* | |
| * application.yml (o minimo que importa): | |
| * | |
| * spring.kafka: | |
| * consumer: | |
| * group-id: ledger-writer | |
| * auto-offset-reset: earliest | |
| * enable-auto-commit: false | |
| * listener: | |
| * ack-mode: RECORD # commit por registro, apos o metodo retornar | |
| */ | |
| @Component | |
| public class TransactionListener { | |
| private static final Logger log = LoggerFactory.getLogger(TransactionListener.class); | |
| private static final String TABLE = "ledger"; | |
| private final DynamoDbClient dynamo; | |
| public TransactionListener(DynamoDbClient dynamo) { | |
| this.dynamo = dynamo; | |
| } | |
| /** | |
| * concurrency = 3 cria 3 containers, cada um com sua propria thread. | |
| * O Kafka distribui particoes distintas entre eles, entao continua valendo | |
| * "uma thread por particao" -- nao ha paralelismo dentro da mesma particao. | |
| * | |
| * O metodo e sincrono de proposito. Anotar com @Async ou devolver | |
| * CompletableFuture faz o container liberar a thread e processar o proximo | |
| * registro antes deste terminar: adeus ordem dentro da particao. | |
| */ | |
| @KafkaListener( | |
| topics = "financial-transactions", | |
| groupId = "ledger-writer", | |
| concurrency = "3") | |
| public void onTransaction(ConsumerRecord<String, TxnEvent> record, Acknowledgment ack) { | |
| TxnEvent event = record.value(); | |
| // ---------- verify ---------- | |
| GetItemResponse existing = dynamo.getItem(GetItemRequest.builder() | |
| .tableName(TABLE) | |
| .key(primaryKey(event)) | |
| // Leitura eventualmente consistente devolveria "nao existe" | |
| // para um item gravado ha poucos milissegundos. | |
| .consistentRead(true) | |
| .build()); | |
| if (existing.hasItem()) { | |
| log.info("Duplicata ignorada: txnId={}", event.txnId()); | |
| ack.acknowledge(); | |
| return; | |
| } | |
| // >>> JANELA DE RACE CONDITION <<< | |
| // Nada garante atomicidade entre o getItem acima e o putItem abaixo. | |
| // Durante um rebalance, o consumer que esta perdendo a particao pode | |
| // estar exatamente aqui enquanto o novo dono ja leu o mesmo offset e | |
| // fez o mesmo getItem. Os dois veem "nao existe" e os dois inserem. | |
| // ---------- insert ---------- | |
| dynamo.putItem(PutItemRequest.builder() | |
| .tableName(TABLE) | |
| .item(toItem(event)) | |
| .build()); | |
| log.info("Gravado: txnId={} account={}", event.txnId(), event.accountId()); | |
| // Offset so avanca depois da escrita -> at-least-once. | |
| // Se o processo morrer antes disto, o registro volta no proximo poll. | |
| ack.acknowledge(); | |
| } | |
| /** | |
| * Mesma responsabilidade, uma chamada so. A condicao e avaliada de forma | |
| * atomica na particao do DynamoDB, entao rebalance, zombie consumer e | |
| * replay deixam de ser problema. | |
| */ | |
| public void onTransactionConditional(TxnEvent event, Acknowledgment ack) { | |
| try { | |
| dynamo.putItem(PutItemRequest.builder() | |
| .tableName(TABLE) | |
| .item(toItem(event)) | |
| .conditionExpression("attribute_not_exists(PK)") | |
| .build()); | |
| log.info("Gravado: txnId={} account={}", event.txnId(), event.accountId()); | |
| } catch (ConditionalCheckFailedException alreadyProcessed) { | |
| // Caminho feliz da idempotencia, nao um erro. | |
| log.info("Duplicata ignorada: txnId={}", event.txnId()); | |
| } | |
| ack.acknowledge(); | |
| } | |
| // ------------------------------------------------------------------ | |
| private Map<String, AttributeValue> primaryKey(TxnEvent event) { | |
| Map<String, AttributeValue> key = new HashMap<>(); | |
| key.put("PK", AttributeValue.fromS("ACCT#" + event.accountId())); | |
| key.put("SK", AttributeValue.fromS("TXN#" + event.txnId())); | |
| return key; | |
| } | |
| private Map<String, AttributeValue> toItem(TxnEvent event) { | |
| Map<String, AttributeValue> item = new HashMap<>(primaryKey(event)); | |
| item.put("amount", AttributeValue.fromN(event.amount())); | |
| item.put("currency", AttributeValue.fromS(event.currency())); | |
| item.put("occurredAt", AttributeValue.fromS(event.occurredAt())); | |
| return item; | |
| } | |
| public record TxnEvent( | |
| String txnId, | |
| String accountId, | |
| String amount, | |
| String currency, | |
| String occurredAt) { | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment