O Aliyun Log Java Producer é uma biblioteca de alto desempenho para gravar logs no SLS a partir de engines de big data, como Flink, Spark e Storm. Ela oferece compressão de logs, envio em lote e transmissão assíncrona, superando as limitações da API ou do SDK padrão.
Pré-requisitos
O Simple Log Service está ativado.
O Simple Log Service SDK for Python está inicializado.
O que é o Aliyun Log Java Producer
O Aliyun Log Java Producer é uma biblioteca de alto desempenho para aplicações Java em cenários de big data e alta concorrência. Em comparação com a API ou o SDK padrão, ele oferece alto desempenho, separação entre lógica de computação e I/O, além de gerenciamento de recursos. A biblioteca aproveita a capacidade de gravação sequencial do SLS para garantir o envio ordenado de logs.
O SLS disponibiliza aplicações de exemplo para ajudar você a começar rapidamente: Aliyun Log Producer Sample Application.
Fluxo de trabalho
Recursos
Segurança de threads: Todos os métodos da interface Producer são seguros para uso concorrente.
Envio assíncrono: O método send retorna imediatamente. O Producer armazena e mescla os dados internamente e os envia em lotes para aumentar o throughput.
Nova tentativa automática: O Producer tenta reenviar falhas com base na contagem de tentativas e no tempo de recuo configurados.
Rastreamento de comportamento: Utilize Callback ou Future para verificar se os dados foram enviados com sucesso e inspecionar cada tentativa de envio para solução de problemas.
Restauração de contexto: Logs originados da mesma instância do Producer compartilham um contexto, o que permite visualizar logs adjacentes no servidor.
Encerramento graceful: O método close processa todos os dados em cache antes de sair e gera as notificações correspondentes.
Benefícios
Em comparação com a API ou o SDK padrão, o Producer oferece:
-
Alto desempenho
Gravar grandes volumes de dados com recursos limitados exige lógica complexa: multithreading, cache, envio em lote e novas tentativas em caso de falha. O Producer gerencia tudo isso, simplifica o desenvolvimento e proporciona vantagens de desempenho.
-
Assíncrono e não bloqueante
Com memória suficiente, o Producer armazena os dados em cache e retorna imediatamente do método send, separando a computação do I/O. Obtenha o resultado por meio do Future retornado ou do Callback fornecido.
-
Recursos controláveis
Configure o limite de memória para dados em cache e o número de threads de envio. Isso permite equilibrar o consumo de recursos e o throughput de gravação.
-
Localização simplificada de problemas
Em caso de falha, o Producer retorna tanto um código de status quanto uma mensagem de erro descritiva. Por exemplo, "connection timeout" para problemas de rede ou "server unresponsive" para servidores que não respondem.
Notas de uso
O aliyun-log-producer chama a operação PutLogs para fazer upload de logs. Limites de tamanho de log bruto se aplicam por solicitação. Para mais informações, consulte Leitura e gravação de dados.
Os recursos do SLS — projetos, Logstores, shards, LogtailConfig, grupos de máquinas, tamanho de LogItem, comprimento de LogItem (Key) e comprimento de LogItem (Value) — estão sujeitos a limites. Para mais informações, consulte Limites básicos de recursos.
Após a primeira execução do código, ative a indexação do Logstore no console do SLS e aguarde um minuto antes de executar consultas.
Ao consultar logs no console, valores de campo que excedem o comprimento máximo são truncados e excluídos da análise. Para mais informações, consulte Criar índices.
Faturamento
Os custos do SDK são iguais aos custos do console. Para mais informações, consulte Visão geral do faturamento.
Etapa 1: Instale o Aliyun Log Java Producer
Para usar o Aliyun Log Java Producer em um projeto Maven, adicione a seguinte dependência ao seu arquivo pom.xml em <dependencies>:
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>aliyun-log-producer</artifactId>
<version>0.3.22</version>
</dependency>
Se ocorrer um conflito de versão, adicione a seguinte dependência em <dependencies>:
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>aliyun-log</artifactId>
<version>0.6.114</version>
<classifier>jar-with-dependencies</classifier>
</dependency>
Etapa 2: Configure o ProducerConfig
O ProducerConfig controla a política de envio. Ajuste os parâmetros conforme sua carga de trabalho:
Config producerConfig = new ProducerConfig();
producerConfig.setTotalSizeInBytes(104857600);
|
Parâmetro |
Tipo |
Descrição |
|
totalSizeInBytes |
Integer |
Tamanho máximo de logs que uma instância do producer pode armazenar em cache. Valor padrão: 100 MB. |
|
maxBlockMs |
Tempo máximo (em segundos) de bloqueio do método send quando o cache do Producer está cheio. Padrão: 60 segundos. Se o tempo limite expirar e o espaço em cache permanecer insuficiente, o método send lança TimeoutException. Defina como 0 para lançar TimeoutException imediatamente quando o espaço em cache for insuficiente. Defina como um valor negativo para bloquear indefinidamente até que haja espaço disponível no cache. |
|
|
ioThreadCount |
Integer |
Número de threads para tarefas de envio. Padrão: número de processadores disponíveis. |
|
batchSizeThresholdInBytes |
Integer |
Limiar de tamanho para enviar um lote. Padrão: 512 KB. Máximo: 5 MB. |
|
batchCountThreshold |
Integer |
Limiar de contagem de logs para enviar um lote. Padrão: 4096. Máximo: 40960. |
|
lingerMs |
Integer |
Tempo de espera antes de enviar um lote. Padrão: 2 segundos. Mínimo: 100 ms. |
|
retries |
Integer |
Contagem máxima de novas tentativas após uma falha inicial de envio. Padrão: 10. Defina como 0 ou menos para enviar o lote diretamente para a fila de falhas na primeira ocorrência. |
|
maxReservedAttempts |
Integer |
Cada tentativa de envio gera um registro de attempt. Este parâmetro controla quantas tentativas recentes são retidas. Padrão: 11. Valores maiores fornecem rastreamento mais detalhado, mas aumentam o uso de memória. |
|
baseRetryBackoffMs |
Integer |
Tempo inicial de recuo para nova tentativa. Padrão: 100 milissegundos. O Producer usa recuo exponencial: tempo de espera antes da N-ésima tentativa = baseRetryBackoffMs × 2^(N-1). |
|
maxRetryBackoffMs |
Integer |
Tempo máximo de recuo para nova tentativa. Padrão: 50 segundos. |
|
adjustShardHash |
Boolean |
Indica se o shardHash deve ser ajustado no envio. Padrão: true. |
|
buckets |
Integer |
Efetivo quando adjustShardHash é true. Reagrupa valores de shardHash no número especificado de buckets para melhorar o envio em lotes. Valores diferentes de shardHash impedem a mesclagem de dados, o que limita o throughput. O reagrupamento permite um envio em lotes mais eficiente. Deve ser uma potência de 2 no intervalo [1, 256]. Padrão: 64. |
Etapa 3: Crie um producer
O Producer suporta autenticação com pares de AccessKey ou tokens STS. Para tokens STS, crie periodicamente um novo ProjectConfig e adicione-o a ProjectConfigs.
LogProducer é a classe de implementação e requer uma instância de ProducerConfig. Crie um Producer da seguinte forma:
Producer producer = new LogProducer(producerConfig);
A criação de um Producer inicia várias threads, o que consome muitos recursos. Compartilhe uma única instância do Producer em toda a aplicação. Todos os métodos do LogProducer são seguros para threads. A tabela abaixo lista as threads internas, onde N é o número da instância começando em 0.
|
Formato do nome da thread |
Quantidade |
Descrição |
|
aliyun-log-producer-<N>-mover |
1 |
Transfere lotes prontos para envio ao pool de threads de envio. |
|
aliyun-log-producer-<N>-io-thread |
ioThreadCount |
Threads no IOThreadPool que executam tarefas de envio de dados. |
|
aliyun-log-producer-<N>-success-batch-handler |
1 |
Processa lotes enviados com sucesso. |
|
aliyun-log-producer-<N>-failure-batch-handler |
1 |
Gerencia lotes com falha no envio. |
Etapa 4: Configure um projeto de log
O ProjectConfig contém o endpoint e as credenciais de acesso para um projeto de destino. Cada projeto requer seu próprio ProjectConfig.
Crie instâncias conforme abaixo:
ProjectConfig project1 = new ProjectConfig("your-project-1", "cn-hangzhou.log.aliyuncs.com", "accessKeyId", "accessKeySecret");
ProjectConfig project2 = new ProjectConfig("your-project-2", "cn-shanghai.log.aliyuncs.com", "accessKeyId", "accessKeySecret");
producer.putProject(project1);
producer.putProject(project2);
Etapa 5: Envie dados
Crie Future ou Callback
Especifique um Callback ao enviar logs. O Callback é invocado na entrega bem-sucedida ou quando ocorre uma exceção durante uma falha de envio.
Se o pós-processamento do resultado for simples e não bloqueante, use o Callback diretamente. Caso contrário, utilize ListenableFuture para lidar com a lógica em um pool de threads separado.
Parâmetros do método:
|
Parâmetro |
Descrição |
|
project |
Projeto de destino para os dados a serem enviados. |
|
logstore |
Logstore de destino para os dados a serem enviados. |
|
logTem |
Dados a serem enviados. |
|
completed |
Tipo atômico Java para garantir que todos os logs sejam enviados (com sucesso ou com falha). |
Enviar dados
A interface Producer fornece vários métodos send com os seguintes parâmetros:
|
Parâmetro |
Descrição |
Obrigatório |
|
project |
Projeto de destino. |
Sim |
|
logStore |
Logstore de destino. |
Sim |
|
logItem |
Logs a serem enviados. |
Sim |
|
topic |
Tópico dos logs. |
Não Nota
Se não especificado, este parâmetro recebe "". |
|
source |
Origem dos logs. |
Não Nota
Se não especificado, este parâmetro recebe o endereço IP do host onde o producer reside. |
|
shardHash |
Valor de hash usado para rotear logs para um shard específico no Logstore. |
Não Nota
Se não especificado, os dados são gravados em um shard aleatório. |
|
callback |
Callback invocado na entrega bem-sucedida ou após todas as tentativas serem esgotadas. |
Não |
Exceções comuns
|
Exceção |
Descrição |
|
TimeoutException |
Lançada quando o tamanho dos logs em cache do Producer excede o limite de memória e não é possível adquirir memória suficiente dentro de maxBlockMs. Se maxBlockMs estiver definido como -1, o bloqueio é indefinido e TimeoutException não ocorre. |
|
IllegalStateException |
Lançada quando o método send é chamado após o fechamento do Producer. |
Etapa 6: Obter o resultado do envio
O método send é assíncrono. Obtenha o resultado por meio do Future retornado ou do Callback fornecido.
Future
O método send retorna um ListenableFuture que suporta registro de callbacks. O exemplo abaixo registra um FutureCallback executado em um pool de threads personalizado. Exemplo completo: SampleProducerWithFuture.java.
package com.aliyun.openservices.aliyun.log.producer.sample;
import com.aliyun.openservices.aliyun.log.producer.*;
import com.aliyun.openservices.aliyun.log.producer.errors.LogSizeTooLargeException;
import com.aliyun.openservices.aliyun.log.producer.errors.MaxBatchCountExceedException;
import com.aliyun.openservices.aliyun.log.producer.errors.ProducerException;
import com.aliyun.openservices.aliyun.log.producer.errors.ResultFailedException;
import com.aliyun.openservices.aliyun.log.producer.errors.TimeoutException;
import com.aliyun.openservices.log.common.LogItem;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class SampleProducerWithFuture {
private static final Logger LOGGER = LoggerFactory.getLogger(SampleProducerWithFuture.class);
private static final ExecutorService EXECUTOR_SERVICE = Executors
.newFixedThreadPool(Math.max(Runtime.getRuntime().availableProcessors(), 1));
public static void main(String[] args) throws InterruptedException {
Producer producer = Utils.createProducer();
int n = 100;
int size = 20;
// The number of logs that have finished (either successfully send, or failed)
final AtomicLong completed = new AtomicLong(0);
for (int i = 0; i < n; ++i) {
List<LogItem> logItems = Utils.generateLogItems(size);
try {
String project = System.getenv("PROJECT");
String logStore = System.getenv("LOG_STORE");
ListenableFuture<Result> f = producer.send(project, logStore, logItems);
Futures.addCallback(
f, new SampleFutureCallback(project, logStore, logItems, completed), EXECUTOR_SERVICE);
} catch (InterruptedException e) {
LOGGER.warn("The current thread has been interrupted during send logs.");
} catch (Exception e) {
if (e instanceof MaxBatchCountExceedException) {
LOGGER.error("The logs exceeds the maximum batch count, e={}", e);
} else if (e instanceof LogSizeTooLargeException) {
LOGGER.error("The size of log is larger than the maximum allowable size, e={}", e);
} else if (e instanceof TimeoutException) {
LOGGER.error("The time taken for allocating memory for the logs has surpassed., e={}", e);
} else {
LOGGER.error("Failed to send logs, e=", e);
}
}
}
Utils.doSomething();
try {
producer.close();
} catch (InterruptedException e) {
LOGGER.warn("The current thread has been interrupted from close.");
} catch (ProducerException e) {
LOGGER.info("Failed to close producer, e=", e);
}
EXECUTOR_SERVICE.shutdown();
while (!EXECUTOR_SERVICE.isTerminated()) {
EXECUTOR_SERVICE.awaitTermination(100, TimeUnit.MILLISECONDS);
}
LOGGER.info("All log complete, completed={}", completed.get());
}
private static final class SampleFutureCallback implements FutureCallback<Result> {
private static final Logger LOGGER = LoggerFactory.getLogger(SampleFutureCallback.class);
private final String project;
private final String logStore;
private final List<LogItem> logItems;
private final AtomicLong completed;
SampleFutureCallback(
String project, String logStore, List<LogItem> logItems, AtomicLong completed) {
this.project = project;
this.logStore = logStore;
this.logItems = logItems;
this.completed = completed;
}
@Override
public void onSuccess(@Nullable Result result) {
LOGGER.info("Send logs successfully.");
completed.getAndIncrement();
}
@Override
public void onFailure(Throwable t) {
if (t instanceof ResultFailedException) {
Result result = ((ResultFailedException) t).getResult();
LOGGER.error(
"Failed to send logs, project={}, logStore={}, result={}", project, logStore, result);
} else {
LOGGER.error("Failed to send log, e=", t);
}
completed.getAndIncrement();
}
}
}
Callback
O Callback executa na thread interna do Producer, e os dados só são liberados após sua conclusão. Evite operações longas no Callback para prevenir bloqueios. Não chame o método send para novas tentativas dentro do Callback — use ListenableFuture em vez disso. Exemplo completo: SampleProducerWithCallback.java.
package com.aliyun.openservices.aliyun.log.producer.sample;
import com.aliyun.openservices.aliyun.log.producer.Callback;
import com.aliyun.openservices.aliyun.log.producer.Producer;
import com.aliyun.openservices.aliyun.log.producer.Result;
import com.aliyun.openservices.aliyun.log.producer.errors.LogSizeTooLargeException;
import com.aliyun.openservices.aliyun.log.producer.errors.ProducerException;
import com.aliyun.openservices.aliyun.log.producer.errors.TimeoutException;
import com.aliyun.openservices.log.common.LogItem;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicLong;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class SampleProducerWithCallback {
private static final Logger LOGGER = LoggerFactory.getLogger(SampleProducerWithCallback.class);
private static final ExecutorService EXECUTOR_SERVICE = Executors.newFixedThreadPool(10);
public static void main(String[] args) throws InterruptedException {
final Producer producer = Utils.createProducer();
int nTask = 100;
// The monotonically increasing sequence number we will put in the data of each log
final AtomicLong sequenceNumber = new AtomicLong(0);
// The number of logs that have finished (either successfully send, or failed)
final AtomicLong completed = new AtomicLong(0);
final CountDownLatch latch = new CountDownLatch(nTask);
for (int i = 0; i < nTask; ++i) {
EXECUTOR_SERVICE.submit(
new Runnable() {
@Override
public void run() {
LogItem logItem = Utils.generateLogItem(sequenceNumber.getAndIncrement());
try {
String project = System.getenv("PROJECT");
String logStore = System.getenv("LOG_STORE");
producer.send(
project,
logStore,
Utils.getTopic(),
Utils.getSource(),
logItem,
new SampleCallback(project, logStore, logItem, completed));
} catch (InterruptedException e) {
LOGGER.warn("The current thread has been interrupted during send logs.");
} catch (Exception e) {
if (e instanceof LogSizeTooLargeException) {
LOGGER.error(
"The size of log is larger than the maximum allowable size, e={}", e);
} else if (e instanceof TimeoutException) {
LOGGER.error(
"The time taken for allocating memory for the logs has surpassed., e={}", e);
} else {
LOGGER.error("Failed to send log, logItem={}, e=", logItem, e);
}
} finally {
latch.countDown();
}
}
});
}
latch.await();
EXECUTOR_SERVICE.shutdown();
Utils.doSomething();
try {
producer.close();
} catch (InterruptedException e) {
LOGGER.warn("The current thread has been interrupted from close.");
} catch (ProducerException e) {
LOGGER.info("Failed to close producer, e=", e);
}
LOGGER.info("All log complete, completed={}", completed.get());
}
private static final class SampleCallback implements Callback {
private static final Logger LOGGER = LoggerFactory.getLogger(SampleCallback.class);
private final String project;
private final String logStore;
private final LogItem logItem;
private final AtomicLong completed;
SampleCallback(String project, String logStore, LogItem logItem, AtomicLong completed) {
this.project = project;
this.logStore = logStore;
this.logItem = logItem;
this.completed = completed;
}
@Override
public void onCompletion(Result result) {
try {
if (result.isSuccessful()) {
LOGGER.info("Send log successfully.");
} else {
LOGGER.error(
"Failed to send log, project={}, logStore={}, logItem={}, result={}",
project,
logStore,
logItem.ToJsonString(),
result);
}
} finally {
completed.getAndIncrement();
}
}
}
}
Etapa 7: Fechar o producer
Feche o Producer quando ele não for mais necessário para garantir que todos os dados em cache sejam processados. Dois modos de encerramento estão disponíveis:
Safe shutdown
Recomendado para a maioria dos casos. O método close() aguarda o processamento de todos os dados em cache, a parada das threads, a execução dos callbacks e a conclusão dos futures antes de retornar.
Este método retorna rapidamente se o callback for não bloqueante. Após o fechamento, os lotes são processados imediatamente sem novas tentativas.
Limited shutdown
Use close(long timeoutMs) para um retorno rápido quando os callbacks puderem bloquear. Se o Producer não for totalmente fechado após timeoutMs, uma IllegalStateException será lançada, indicando possível perda de dados e callbacks não executados.
Perguntas frequentes
Existem limitações no número de operações de gravação de dados?
As operações de leitura e gravação no SLS estão sujeitas a limites de tamanho. Para mais informações, consulte Leitura e gravação de dados.
Os recursos do SLS — projetos, Logstores, shards, LogtailConfig, grupos de máquinas, tamanho de LogItem, comprimento de LogItem (Key) e comprimento de LogItem (Value) — estão sujeitos a limites. Para mais informações, consulte Limites básicos de recursos.
O que fazer se nenhum dado for gravado no SLS?
Se nenhum dado for gravado no SLS, solucione o problema da seguinte forma:
Verifique se as versões dos JARs
aliyun-log-producer,aliyun-logeprotobuf-javacorrespondem à documentação de instalação. Atualize-as se necessário.O método send é assíncrono. Use um Callback ou Future para determinar a causa da falha.
Se o método onCompletion do Callback não for chamado, garanta que
producer.close()seja invocado antes da saída do programa. Chamarproducer.close()libera todos os dados em cache.O Producer usa SLF4J para logging. Configure um framework de logging e ative o nível DEBUG para verificar erros.
Se o problema persistir, abra um ticket.
Referências
Além de seu SDK nativo, o SLS também suporta os SDKs comuns da Alibaba Cloud. Para mais informações, consulte Simple Log Service_SDK Center_Alibaba Cloud OpenAPI Explorer.