Skip to content

Instantly share code, notes, and snippets.

@rponte
Created September 18, 2026 13:18
Show Gist options
  • Select an option

  • Save rponte/27dd0b4bb03b4a13cc910f5aec555b68 to your computer and use it in GitHub Desktop.

Select an option

Save rponte/27dd0b4bb03b4a13cc910f5aec555b68 to your computer and use it in GitHub Desktop.
A naive idempotent Kafka consumer written in Java and Spring Boot
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