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 segundos.
Visão geral
Um Logstore contém múltiplos shards. Um consumer group consome dados atribuindo esses shards aos seus consumidores conforme as regras a seguir:
Cada shard é atribuído a apenas um consumidor por vez.
Um único consumidor pode receber vários shards.
Quando um novo consumidor entra em um consumer group, os shards são rebalanceados entre todos os consumidores para garantir o equilíbrio de carga. As mesmas regras se aplicam.
Conceitos principais
|
Termo |
Descrição |
|
consumer group |
Conjunto de consumidores que consomem dados conjuntamente do mesmo Logstore, sem duplicação. Importante
É possível criar até 30 consumer groups por Logstore. |
|
consumer |
Unidade básica de um consumer group responsável pelo consumo efetivo dos dados. Importante
Os consumidores no mesmo consumer group devem ter nomes exclusivos. |
|
Logstore |
Unidade de 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 ficam sempre armazenados em um shard. Para mais informações, consulte shard. |
|
checkpoint |
Posição no fluxo de dados que marca o último dado processado por um consumidor. Permite retomar 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 consumo continua a partir do último checkpoint, evitando duplicidade. |
Pré-requisitos
Um projeto e um Logstore padrão foram criados e os logs estão sendo coletados. Para mais informações, consulte Manage projects, Create a basic Logstore e Data collection.
Você configurou as credenciais de acesso.
Você inicializou o SLS SDK for Java.
Etapa 1: Criar um consumer group
Não é possível criar um consumer group diretamente no console web do SLS. Utilize o SDK, a API ou a CLI para essa finalidade.
Crie um consumer group usando o SDK, a API ou a CLI.
SDK
O código a seguir cria um consumer group:
Para exemplos de código sobre gerenciamento de consumer groups, consulte Use the Java SDK to manage consumer groups e Use Simple Log Service SDK for Python to manage consumer groups.
API
Para criar um consumer group via API, consulte CreateConsumerGroup.
Para verificar se o consumer group foi criado, consulte ListConsumerGroup.
CLI
Para criar um consumer group via CLI, consulte create_consumer_group.
Para verificar se o consumer group foi criado, consulte list_consumer_group.
Etapa 2: Consumir dados de log
Como funciona
Na primeira execução de um consumidor com o SDK de consumer group, o SDK cria o grupo caso ele ainda não exista. O checkpoint inicial define a posição de início do consumo e aplica-se apenas na criação do consumer group. Em reinicializações posteriores, o consumidor retoma a partir do último checkpoint salvo no servidor. Exemplo:
LogHubConfig.ConsumePosition.BEGIN_CURSOR: O consumer group inicia o consumo a partir do primeiro log no Logstore.LogHubConfig.ConsumePosition.END_CURSOR: O consumer group inicia o consumo após o último log no Logstore.
Exemplos
É possível consumir dados com um consumer group usando os SDKs para Java, C++, Python e Go. Os exemplos a seguir usam Java.
Exemplo 1: Consumo via SDK
-
Adicione as dependências do Maven.
-
Crie uma classe para definir a lógica de processamento de 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 processador de logs.
-
Crie uma classe principal para configurar e iniciar a thread do 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 de logs.
-
Crie uma factory para gerar instâncias do processador de logs.
-
Crie uma classe principal para configurar e iniciar a thread do worker.
-
Execute o arquivo
Main.java.
Etapa 3: Visualizar o status do consumer group
Visualize o status de um consumer group por um dos métodos a seguir:
Java SDK
-
Verifique o checkpoint de consumo de cada shard. O código a seguir serve como exemplo:
-
Console
Faça login no console do Simple Log Service.
-
Na seção Projects, clique no projeto desejado.

Na aba , clique no ícone
à esquerda do Logstore desejado e, em seguida, clique no ícone
à esquerda de Data Consumption.Na lista de consumer groups, clique no grupo desejado.
Na página Consumer Group Status, visualize o checkpoint de consumo de cada shard. Esta página exibe detalhes de cada shard, incluindo ID (shard), Last Consumed Time e o Client responsável pelo consumo. Os botões Refresh e Reset Checkpoint também estão disponíveis.
Operações relacionadas
-
Autorização de usuário RAM
Para gerenciar consumer groups com um usuário RAM, conceda as permissões necessárias. Para mais informações, consulte Create and authorize a RAM user.
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 para 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
Region ID: cn-hangzhou
Nome do projeto: project-test
Nome do Logstore: logstore-test
Nome do consumer group: consumergroup-test
Alternativamente, anexe a política de sistema
AliyunLogFullAccessao usuário RAM. Essa política concede permissões completas para gerenciar o Simple Log Service e já inclui todas as permissões exigidas pelas operações de consumer group, eliminando a necessidade de criar uma política personalizada. -
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. Veja abaixo 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 conforme suas necessidades.
Se já existir um checkpoint salvo no servidor, o consumo será retomado a partir dele.
O Simple Log Service prioriza um checkpoint salvo para o consumo. Caso especifique um horário de início, garanta 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