O consumo de dados do Log Service em tempo real via SDK exige o gerenciamento de detalhes de implementação, como balanceamento de carga e failover entre consumidores. Um consumer group gerencia essas complexidades automaticamente, permitindo o consumo de dados quase em tempo real, geralmente em questão de segundos.
Visão geral
Um Logstore contém múltiplos shards. Um consumer group consome dados atribuindo esses shards aos seus consumidores conforme as seguintes regras:
Cada shard é atribuído a apenas um consumidor por vez.
Um único consumidor pode receber múltiplos shards.
Quando um novo consumidor entra em um consumer group, ocorre um rebalanceamento dos shards entre todos os consumidores para garantir o equilíbrio de carga. As mesmas regras se aplicam.
Conceitos principais
|
Termo |
Descrição |
|
consumer group |
Conjunto formado por vários consumidores que consomem dados do mesmo Logstore simultaneamente, sem duplicidade. Importante
É possível criar até 30 consumer groups para cada Logstore. |
|
consumer |
Unidade básica de um consumer group responsável pelo consumo efetivo de dados. Importante
Os consumidores dentro do mesmo consumer group devem ter nomes exclusivos. |
|
Logstore |
Unidade destinada à coleta, armazenamento e consulta de dados. Para mais informações, consulte Logstore. |
|
shard |
Unidade que controla a capacidade de leitura e gravação de um Logstore. Os dados são sempre armazenados em um shard. Para mais informações, consulte shard. |
|
checkpoint |
Posição no fluxo de dados que marca o último dado processado pelo consumidor. Isso permite que o consumidor retome o processamento desse ponto após uma reinicialização. Nota
Ao consumir dados com um consumer group, os checkpoints são salvos automaticamente em caso de falha do programa. Após a recuperação, o programa continua o consumo a partir do último checkpoint, evitando o consumo duplicado. |
Pré-requisitos
Um projeto e um Logstore padrão foram criados, e os logs já estão sendo coletados. Para mais informações, consulte Gerenciar projetos, Criar um Logstore básico e Coleta de dados.
Etapa 1: Criar um consumer group
É possível criar um consumer group utilizando o SDK, a API ou a CLI.
SDK
O código abaixo cria um consumer group:
Para exemplos de código sobre como gerenciar consumer groups, consulte Usar o Java SDK para gerenciar consumer groups e Usar o Simple Log Service SDK for Python para gerenciar consumer groups.
API
Para criar um consumer group usando a API, consulte CreateConsumerGroup.
Para verificar se o consumer group foi criado, consulte ListConsumerGroup.
CLI
Para criar um consumer group usando a CLI, consulte create_consumer_group.
Para verificar se o consumer group foi criado, consulte list_consumer_group.
Etapa 2: Consumir dados de log
Funcionamento
Quando um consumidor que utiliza o SDK de consumer group é iniciado pela primeira vez, o SDK cria o consumer group caso ele ainda não exista. O checkpoint inicial define a posição de começo do consumo, sendo utilizado apenas na criação do consumer group. Nas reinicializações seguintes, o consumidor retoma a partir do último checkpoint salvo no servidor. Por exemplo:
LogHubConfig.ConsumePosition.BEGIN_CURSOR: O consumer group inicia o consumo a partir do primeiro log existente no Logstore.LogHubConfig.ConsumePosition.END_CURSOR: O consumer group inicia o consumo após o último log existente no Logstore.
Exemplos
É possível consumir dados com um consumer group usando os SDKs para Java, C++, Python e Go. Os exemplos a seguir utilizam Java.
Exemplo 1: Consumo via SDK
-
Adicione as dependências do Maven.
-
Crie uma classe para definir a lógica de processamento dos logs.
Para mais exemplos de código, consulte os repositórios aliyun-log-consumer-java e Aliyun LOG Go Consumer.
-
Crie uma factory para gerar instâncias do seu processador de logs.
-
Crie uma classe principal para configurar e iniciar a thread worker.
-
Execute o arquivo
Main.java.
Exemplo 2: Consumo via SDK com SPL
-
Adicione as dependências do Maven.
-
Crie uma classe para definir a lógica de processamento dos logs.
-
Crie uma factory para gerar instâncias do seu processador de logs.
-
Crie uma classe principal para configurar e iniciar a thread worker.
-
Execute o arquivo
Main.java.
Etapa 3: Visualizar o status do consumer group
Visualize o status de um consumer group usando um dos métodos a seguir:
Java SDK
-
Visualize o checkpoint de consumo de cada shard. O código abaixo serve como exemplo:
-
Console
Faça login no console do Simple Log Service.
-
Na seção Projects, clique em no projeto desejado.

Na aba , clique em no ícone
à esquerda do Logstore alvo e, em seguida, clique em no ícone
à esquerda de Data Consumption.Na lista de consumer groups, clique em no grupo desejado.
Na página Consumer Group Status, visualize o checkpoint de consumo de cada shard. Esta página exibe detalhes individuais de cada shard, incluindo seu ID (shard), Last Consumed Time e o Client responsável pelo consumo. Também estão disponíveis os botões Refresh e Reset Checkpoint.
Operações relacionadas
-
Autorização de usuário RAM
Para gerenciar consumer groups com um usuário RAM, conceda as permissões necessárias a esse usuário. Para mais informações, consulte Criar e autorizar um usuário RAM.
A tabela a seguir lista as Actions necessárias.
Action
Descrição
Resource
log:GetCursorOrData(GetCursor)
Obtém um cursor com base em um horário especificado.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}
log:CreateConsumerGroup(CreateConsumerGroup)
Cria um consumer group em um Logstore específico.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
log:ListConsumerGroup(ListConsumerGroup)
Lista todos os consumer groups em um Logstore específico.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/*
log:ConsumerGroupUpdateCheckPoint(UpdateCheckPoint)
Atualiza o checkpoint em um shard para um consumer group específico.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
log:ConsumerGroupHeartBeat(ConsumerGroupHeartBeat)
Envia um heartbeat de um consumidor específico para o servidor.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
log:UpdateConsumerGroup(UpdateConsumerGroup)
Modifica as propriedades de um consumer group específico.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
log:GetConsumerGroupCheckPoint(GetConsumerGroupCheckPoint)
Obtém o checkpoint de um ou de todos os shards de um consumer group específico.
acs:log:${regionName}:${projectOwnerAliUid}:project/${projectName}/logstore/${LogStoreName}/consumergroup/${consumerGroupName}
Para conceder as permissões listadas abaixo a um usuário RAM, utilize a política de exemplo a seguir.
ID da conta Alibaba Cloud: 174649****602745
ID da região: cn-hangzhou
Nome do projeto: project-test
Nome do Logstore: logstore-test
Nome do consumer group: consumergroup-test
-
Solução de problemas
Para facilitar a solução de problemas, configure o Log4j na sua aplicação consumidora para registrar exceções do consumer group. Abaixo está um exemplo de configuração do log4j.properties:
log4j.rootLogger = info,stdout log4j.appender.stdout = org.apache.log4j.ConsoleAppender log4j.appender.stdout.Target = System.out log4j.appender.stdout.layout = org.apache.log4j.PatternLayout log4j.appender.stdout.layout.ConversionPattern = [%-5p] %d{yyyy-MM-dd HH:mm:ss,SSS} method:%l%n%m%nApós configurar o Log4j, a aplicação consumidora emitirá informações de exceção semelhantes às seguintes:
[WARN ] 2018-03-14 12:01:52,747 method:com.aliyun.openservices.loghub.client.LogHubConsumer.sampleLogError(LogHubConsumer.java:159) com.aliyun.openservices.log.exception.LogException: Invalid loggroup count, (0,1000] -
Consumir dados a partir de um horário específico
// consumerStartTimeInSeconds indicates the time from which to start consuming data. public LogHubConfig(String consumerGroupName, String consumerName, String loghubEndPoint, String project, String logStore, String accessId, String accessKey, int consumerStartTimeInSeconds); // position is an enumeration. LogHubConfig.ConsumePosition.BEGIN_CURSOR starts consumption from the oldest data. LogHubConfig.ConsumePosition.END_CURSOR starts consumption from the latest data. public LogHubConfig(String consumerGroupName, String consumerName, String loghubEndPoint, String project, String logStore, String accessId, String accessKey, ConsumePosition position);NotaEscolha um construtor de acordo com suas necessidades.
Caso já exista um checkpoint salvo no servidor, o consumo será retomado a partir desse ponto.
O Simple Log Service prioriza um checkpoint salvo para o consumo. Se você especificar um horário de início, certifique-se de que o valor de consumerStartTimeInSeconds esteja dentro do período de retenção de dados (TTL). Caso contrário, o horário de início especificado será ignorado.
-
Redefinir um checkpoint
public static void updateCheckpoint() throws Exception { Client client = new Client(host, accessId, accessKey); // The timestamp must be a UNIX timestamp in seconds. If your timestamp is in milliseconds, divide it by 1000. long timestamp = Timestamp.valueOf("2017-11-15 00:00:00").getTime() / 1000; ListShardResponse response = client.ListShard(new ListShardRequest(project, logStore)); for (Shard shard : response.GetShards()) { int shardId = shard.GetShardId(); String cursor = client.GetCursor(project, logStore, shardId, timestamp).GetCursor(); client.UpdateCheckPoint(project, logStore, consumerGroup, shardId, cursor); } }
Referências
-
API
Actions
API
Criar um consumer group
Consultar um consumer group
Excluir um consumer group
Atualizar um consumer group
Enviar heartbeat do consumidor
Consultar checkpoints do consumer group
Atualizar checkpoints do consumer group
-
SDK
Linguagem
Referências
Java
Python
-
CLI
Actions
CLI
Criar um consumer group
Consultar um consumer group
Atualizar um consumer group
Excluir um consumer group
Consultar checkpoints do consumer group
Atualizar checkpoints do consumer group