Quando os consumidores ficam atrasados em relação aos produtores, as mensagens se acumulam e a latência de processamento aumenta drasticamente. Offsets mal configurados podem causar o salto ou reprocessamento silencioso de mensagens. Este guia aborda padrões comprovados para criar aplicações consumidoras confiáveis e de alto throughput no ApsaraMQ for Kafka, incluindo grupos de consumidores, gerenciamento de offsets, tratamento de falhas e ajuste de desempenho.
Como funciona o consumo de mensagens
Um consumidor do ApsaraMQ for Kafka repete um ciclo de três etapas:
Poll -- Busca um lote de mensagens do broker.
Processamento -- Executa a lógica de negócios nas mensagens buscadas.
Novo poll -- Busca o próximo lote após a conclusão do processamento.
Mantenha o tempo de processamento curto para preservar um intervalo de poll constante. Atrasos prolongados no processamento disparam rebalanceamentos ou provocam acúmulo de mensagens.
Grupos de consumidores e balanceamento de carga
Um grupo de consumidores é um conjunto de instâncias de consumidor que compartilham o mesmo group.id. O ApsaraMQ for Kafka distribui as partições dos tópicos assinados uniformemente entre os consumidores de um grupo, garantindo que cada mensagem seja entregue a exatamente um consumidor.
Exemplo: O Grupo A assina o Tópico A. Três consumidores — C1, C2 e C3 — pertencem ao grupo. Cada mensagem recebida vai para C1, C2 ou C3, equilibrando a carga de consumo entre os três.
O ApsaraMQ for Kafka aciona um rebalanceamento quando os consumidores são:
Iniciados pela primeira vez
Reiniciados
Adicionados ou removidos do grupo
Evite rebalanceamentos frequentes
Os rebalanceamentos pausam o consumo enquanto as partições são reatribuídas. Rebalanceamentos frequentes prejudicam o throughput e aumentam a latência de processamento. Eles ocorrem geralmente quando os heartbeats dos consumidores expiram.
Para evitar esse problema, ajuste os parâmetros relevantes ou aumente a taxa de consumo. Para obter etapas detalhadas de solução de problemas, consulte Por que ocorrem rebalanceamentos frequentes no meu cliente consumidor?
Configure as partições
A contagem de partições determina quantos consumidores podem processar mensagens em paralelo. Cada partição é atribuída a exatamente um consumidor dentro de um grupo. Se houver mais consumidores do que partições, os excedentes ficarão ociosos.
Contagem recomendada de partições
A contagem padrão de partições no console do ApsaraMQ for Kafka é 12, valor adequado para a maioria das cargas de trabalho. Ao dimensionar o sistema, siga estas diretrizes:
|
Contagem de partições |
Impacto |
|
Menor que 12 |
Pode reduzir o desempenho de produção e consumo de mensagens |
|
12 -- 100 |
Faixa recomendada para a maioria das cargas de trabalho |
|
Maior que 100 |
Pode desencadear rebalanceamentos frequentes de consumidores |
Não é possível reduzir a contagem de partições após um aumento. Incremente as partições gradualmente, em vez de fazer grandes saltos de uma só vez.
Padrões de assinatura
O ApsaraMQ for Kafka oferece suporte a dois padrões de assinatura.
Um grupo assina vários tópicos
Um único grupo de consumidores pode assinar múltiplos tópicos. As mensagens de todos os tópicos assinados são distribuídas uniformemente entre os consumidores do grupo.
String topicStr = kafkaProperties.getProperty("topic");
String[] topics = topicStr.split(",");
for (String topic: topics) {
subscribedTopics.add(topic.trim());
}
consumer.subscribe(subscribedTopics);
Vários grupos assinam um tópico
Diversos grupos de consumidores podem assinar o mesmo tópico de forma independente. Cada grupo recebe todas as mensagens e opera sem interferir nos demais.
Utilize este padrão quando diferentes aplicações precisarem processar o mesmo fluxo de dados independentemente — por exemplo, um grupo para análises em tempo real e outro para arquivamento de dados.
Mantenha um grupo por aplicação
Use um grupo de consumidores dedicado para cada aplicação. Caso uma única aplicação precise executar lógicas de consumo distintas, crie arquivos de configuração separados (por exemplo, kafka1.properties e kafka2.properties) com valores diferentes para group.id.
Alinhe as assinaturas dentro de um grupo
Todos os consumidores do mesmo grupo devem assinar o mesmo conjunto de tópicos. Assinaturas divergentes dificultam a solução de problemas e podem resultar em comportamentos inesperados de rebalanceamento.
Gerencie os offsets do consumidor
Cada partição rastreia um offset máximo, que representa o número total de mensagens recebidas. Cada consumidor rastreia um offset do consumidor, indicando quantas mensagens ele já processou naquela partição. A diferença entre esses dois valores corresponde ao acúmulo de mensagens (backlog não consumido).
Commit automático versus manual
O ApsaraMQ for Kafka disponibiliza dois parâmetros para confirmar os offsets do consumidor:
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Ativa a confirmação automática de offsets |
|
|
|
Intervalo entre as confirmações automáticas |
Com o commit automático ativado, o cliente verifica o tempo decorrido desde a última confirmação antes de cada poll. Se o tempo ultrapassar o valor definido em auto.commit.interval.ms, ele confirma o offset atual.
Ao usar o commit automático, garanta que todas as mensagens do poll anterior sejam totalmente processadas antes do próximo poll. Se o consumidor gravar em um armazenamento externo e essa gravação falhar, o commit automático ainda poderá avançar o offset. Isso faz com que essas mensagens sejam ignoradas permanentemente em vez de serem tentadas novamente. Compreenda essa contrapartida antes de depender do comportamento padrão.
Confirme offsets manualmente
Para confirmar offsets manualmente, defina enable.auto.commit como false e chame a função commit(offsets) após processar cada lote.
Comportamento de redefinição de offset
Os offsets do consumidor são redefinidos em duas situações:
Nenhum offset confirmado existe — por exemplo, quando um consumidor se conecta a um broker pela primeira vez.
O offset confirmado é inválido — por exemplo, o offset máximo em uma partição é 10, mas o consumidor tenta ler a partir do offset 11.
Controle o comportamento de redefinição em clientes Java com auto.offset.reset:
|
Valor |
Comportamento |
|
|
Inicia a partir da mensagem mais recente. Escolha esta opção como padrão |
|
|
Inicia a partir da mensagem disponível mais antiga |
|
|
Lança uma exceção em vez de redefinir |
Prefira latest em vez de earliest para evitar o reprocessamento de todo o histórico do tópico quando ocorrer um offset inválido. Se você confirmar offsets manualmente com a função commit(offsets), none é uma escolha segura, pois seus offsets devem estar sempre válidos.
Lide com mensagens grandes
Quando mensagens individuais excedem os tamanhos típicos, ajuste estes parâmetros para controlar o comportamento de busca:
|
Parâmetro |
Orientação |
|
|
Número máximo de registros por poll. Defina como |
|
|
Quantidade máxima de dados por solicitação de fetch. Defina um valor ligeiramente superior ao tamanho esperado da mensagem |
|
|
Limite de dados por partição em cada fetch. Configure um pouco acima do tamanho previsto da mensagem |
Com essas configurações, o consumidor busca mensagens grandes uma de cada vez, evitando pressão na memória causada por lotes excessivamente grandes.
Trate a duplicação de mensagens
O ApsaraMQ for Kafka utiliza semântica de entrega at-least-once: cada mensagem é entregue pelo menos uma vez para evitar perda de dados, mas duplicações podem ocorrer durante erros de rede ou reinicializações do cliente.
Implemente consumo idempotente
Se sua aplicação for sensível a duplicatas — como transações financeiras ou processamento de pedidos — implemente a deduplicação no nível da aplicação:
Anexe uma chave única a cada mensagem durante a produção (por exemplo, um ID de transação).
Antes de processar uma mensagem consumida, verifique se a chave já foi processada anteriormente.
Ignore mensagens cujas chaves já tenham sido processadas.
Caso sua aplicação tolere duplicatas ocasionais (como agregação de métricas ou ingestão de logs), pule esta etapa.
Trate falhas de consumo
As mensagens dentro de uma partição são consumidas sequencialmente. Quando o processamento falha — devido a dados malformados ou indisponibilidade de um serviço downstream, por exemplo — escolha uma destas estratégias:
|
Estratégia |
Contrapartida |
|
Tentativa no local |
Tenta reprocessar a mensagem com falha até obter sucesso. Simples de implementar, mas bloqueia a thread do consumidor e causa acúmulo de mensagens naquela partição |
|
Tópico de dead-letter |
Envia mensagens com falha para um tópico dedicado para inspeção posterior. O consumo continua sem bloqueios, mas exige um processo separado para investigar e reprocessar as falhas |
Para a maioria das cargas de trabalho em produção, a abordagem de dead-letter evita o travamento do pipeline do consumidor.
Reduza a latência de consumo
O ApsaraMQ for Kafka usa um modelo baseado em pull: os consumidores buscam mensagens do broker em cada ciclo de poll. Quando os consumidores acompanham o ritmo dos produtores, a latência permanece baixa. Um pico de latência geralmente indica acúmulo de mensagens.
Causas comuns de acúmulo
|
Causa |
Solução |
|
Taxa de consumo inferior à taxa de produção |
Adicione consumidores ou aumente o throughput de processamento |
|
Thread do consumidor bloqueada por chamada remota lenta |
Defina timeouts nas chamadas externas para liberar a thread após uma espera limitada |
Aumente o throughput de consumo
Adicione mais consumidores. Inicie instâncias adicionais de consumidores (uma thread por consumidor) ou implante mais processos de consumo. O número de consumidores ativos não pode exceder a contagem de partições; os excedentes permanecerão ociosos.
Utilize um pool de threads para processamento. Desacople o polling do processamento para paralelizar o trabalho:
Defina um pool de threads com um número fixo de worker threads.
Faça o poll de um lote de mensagens.
Envie cada mensagem para o pool de threads para processamento concorrente.
Após a conclusão de todas as tarefas do lote, faça o poll do próximo lote.
Esse padrão mantém o loop de poll responsivo enquanto distribui o trabalho de processamento entre várias threads.
Filtre mensagens
O ApsaraMQ for Kafka não fornece filtragem nativa de mensagens. Implemente a filtragem no nível da aplicação usando uma das seguintes abordagens:
|
Abordagem |
Quando usar |
|
Tópicos separados |
Pequeno número de categorias distintas de mensagens |
|
Filtragem no lado do cliente |
Muitas categorias ou lógica de filtragem que muda frequentemente |
Combine ambas as abordagens quando seu caso de uso exigir: direcione categorias amplas para tópicos diferentes e aplique filtros refinados no lado do consumidor.
Transmita mensagens em broadcast
O ApsaraMQ for Kafka não suporta nativamente o broadcast de mensagens. Para entregar cada mensagem a múltiplos consumidores independentes, crie um grupo de consumidores separado para cada consumidor que precisar receber todas as mensagens. Cada grupo consome independentemente o fluxo completo de mensagens do tópico.