AdminClient em Kafka: Arquitetura e Aplicações Práticas

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 Call representando cada operação.
  • Traduz solicitações (ex: CreateTopicsRequest) em requisições serializáveis.
  • Envia os Call para 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 ArrayList e HashMap em 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.

Tags: kafka adminclient java Distributed-Systems monitoring

Publicado em 9-26 11:12