Abordagem Moderna para Gerenciamento de Kafka com AdminClient
A partir da versão 0.11, o Apache Kafka introduziu uma nova abordagem para operações administrativas por meio do AdminClient, um cliente Java voltado para tarefas de gerenciamento. Este componente substituiu scripts CLI tradicionais e ferramentas internas baseadas em ZooKeeper, oferecendo maiorr segurança, escalabilidade e integração com aplicações.
Limitações dos Scripts Clássicos
Os comandos como kafka-topics.sh apresentam problemas críticos:
- Executam apenas em ambiente terminal, dificultando integração com sistemas automatizados.
- Dependem diretamente do ZooKeeper, evitando mecanismos de autenticação do Kafka (como SASL/SCRAM), permitindo criação de tópicos mesmo por usuários sem permissão.
- Utilizam classes internas do servidor Kafka, o que vai contra a filosofia de separação entre cliente e servidor.
O novo AdminClient reside no pacote org.apache.kafka.clients.admin e opera via protocolo Kafka, respeitando políticas de segurança e autorização definidas no cluster.
Dependência e Configuração
Para usar o AdminClient, inclua a dependência no projeto:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.3.0</version>
</dependency>
Ou com Gradle:
implementation 'org.apache.kafka:kafka-clients:2.3.0'
Funcionalidades Principais
O AdminClient suporta operações essenciais:
- Criação, exclusão e listagem de tópicos
- Gestão de permissões por usuário ou grupo
- Consulta e alteração de configurações globais e específicas (Broker, tópico, cliente)
- Monitoramento de logs de réplicas (path, tamanho, uso)
- Aumento de partições (repartição dinâmica)
- Remoção de mensagens anteriores a uma posição específica
- Geração e gestão de tokens de delegação (Delegation Tokens)
- Controle de grupos de consumidores (listagem, deslocamento, limpeza)
- Eleição de líderes preferenciais para partições
Modelo de Execução Multithreaded
O AdminClient utiliza um modelo de dois threads para garantir eficiência e isolamento:
Thread Principal (Frontend)
- Cria objetos
Callrepresentando cada operação. - Traduz solicitações (ex:
CreateTopicsRequest) em requisições serializáveis. - Envia os
Callpara uma fila de novas requisições. - Aguarda resultados via
Future.
Thread I/O (Backend)
- Trabalha com três estruturas internas:
- Nova Requisição: Recebe chamadas do frontend.
- Requisições Pendentes: Preparadas para envio ao Broker.
- Requisições Ativas: Em processamento ativo.
Essa divisão permite que o frontend não fique bloqueado por locks, enquanto o backend processa requisições com segurança usando monitor e notificações via wait/notify.
Observação: O uso de
ArrayListeHashMapem vez de esrtuturas padrão de fila é intencional — otimizações de desempenho e controle granular sobre o ciclo de vida das requisições.
Criação e Encerramento do Cliente
É obrigatório criar e fechar o AdminClient corretamente:
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "broker-host:9092");
props.put("request.timeout.ms", "600000");
try (AdminClient client = AdminClient.create(props)) {
// Operações aqui
}
O método close() deve ser chamado explicitamente, mesmo com try-with-resources, para liberar recursos de conexão e thread.
Exemplos de Uso
Criar um Tópico
String topicName = "demo-topic";
NewTopic topic = new NewTopic(topicName, 5, (short) 2);
try (AdminClient client = AdminClient.create(props)) {
CreateTopicsResult result = client.createTopics(Arrays.asList(topic));
result.all().get(10, TimeUnit.SECONDS); // Espera conclusão
}
O NewTopic define nome, número de partições e fator de replicação.
Consultar Deslocamentos de Grupo de Consumidores
String groupId = "consumer-group-1";
try (AdminClient client = AdminClient.create(props)) {
ListConsumerGroupOffsetsResult result = client.listConsumerGroupOffsets(groupId);
Map<TopicPartition, OffsetAndMetadata> offsets =
result.partitionsToOffsetAndMetadata().get(10, TimeUnit.SECONDS);
offsets.forEach((tp, meta) ->
System.out.printf("Tópico %s, Partição %d: %d%n",
tp.topic(), tp.partition(), meta.offset()));
}
Este exemplo extrai o deslocamento atual de cada partição no grupo.
Medir Uso de Disco por Broker
Como métricas JMX não fornecem esse dado diretamente, o AdminClient permite calcular o uso:
int brokerId = 1;
try (AdminClient client = AdminClient.create(props)) {
DescribeLogDirsResult result = client.describeLogDirs(Collections.singletonList(brokerId));
long totalSize = result.all().get().values().stream()
.flatMap(map -> map.values().stream())
.flatMap(logDirInfo -> logDirInfo.replicaInfos.values().stream())
.mapToLong(replicaInfo -> replicaInfo.size)
.sum();
System.out.println("Uso total de disco no broker " + brokerId + ": " + totalSize + " bytes");
}
Esse cálculo acumula o tamanho dos arquivos de log de todas as partições replicadas em um Broker específico.
Diagnóstico e Problemas Comuns
Se operações não retornam resultados ou travam, verifique se a thread I/O está funcionando. Ela pode parar inesperadamente por exceções não tratadas.
Use jstack para analisar o estado da thread com prefixo kafka-admin-client-thread. Um deadlock nessa thread geralmente indica bug na implementação — embora o cliente funcione normalmente no frontend, as operações falham silenciosamente.