O SLS permite a ingestão de logs pelo protocolo Kafka a partir de coletores como Kafka Producer SDK, Beats, Collectd, Fluentd, Logstash, Telegraf, Fluent-bit e Vector.
Limites
A versão mínima suportada do protocolo Kafka é a 2.1.0.
O protocolo de conexão SASL_SSL é obrigatório para garantir a transmissão segura dos logs.
Permissões
Sua conta deve ter uma das seguintes permissões:
-
Esta política concede permissões para gerenciar o SLS. Para obter instruções sobre autorização, consulte Gerenciar permissões de usuário RAM e Gerenciar permissões para uma função RAM.
-
Políticas personalizadas
-
Crie uma política personalizada. Na aba Script Editor, substitua o conteúdo existente pelo script abaixo. Para mais informações, consulte Criar uma política personalizada.
NotaNo script, substitua
project_namepelo nome real do seu projeto.{ "Version": "1", "Statement": [ { "Action": "log:GetProject", "Resource": "acs:log:*:*:project/project_name", "Effect": "Allow" }, { "Action": [ "log:GetLogStore", "log:ListShards", "log:PostLogStoreLogs" ], "Resource": "acs:log:*:*:project/project_name/logstore/*", "Effect": "Allow" } ] } Anexe a política personalizada a um usuário RAM. Para mais informações, consulte Gerenciar permissões de usuário RAM.
-
Parâmetros de configuração
Configure os parâmetros a seguir para enviar logs pelo protocolo Kafka.
|
Nome da configuração |
Valor da configuração |
Descrição |
Exemplo |
|
SLS_KAFKA_ENDPOINT |
Endpoint de conexão inicial no formato |
|
aliyun-project-test é o nome do projeto.
|
|
SLS_PROJECT |
Nome do projeto |
Nome do projeto SLS. |
aliyun-project-test |
|
SLS_LOGSTORE |
Nome do Logstore |
Nome do Logstore. Se você adicionar |
Por exemplo, se o nome do Logstore for
|
|
SLS_PASSWORD |
AccessKey secret com permissões de gravação no SLS. |
Para saber o que é um par de AccessKeys e como criar um, consulte Criar um par de AccessKeys. O valor consiste no AccessKey ID e no AccessKey secret, separados pelo símbolo
|
LTAI#yourAccessKeySecret |
Para usar um grupo de consumidores Kafka e consumir dados do SLS em tempo real, envie um ticket para entrar em contato com o suporte técnico da Alibaba Cloud.
Exemplo 1: Enviar logs com Beats
Os coletores da família Beats (MetricBeat, PacketBeat, Winlogbeat, Auditbeat, Filebeat, Heartbeat) podem enviar logs para o SLS usando o protocolo Kafka.
-
Exemplo de configuração
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.output.kafka: # initial brokers for reading cluster metadata hosts: ["SLS_KAFKA_ENDPOINT"] username: "SLS_PROJECT" password: "SLS_PASSWORD" ssl.certificate_authorities: # message topic selection + partitioning topic: 'SLS_LOGSTORE' partition.round_robin: reachable_only: false required_acks: 1 compression: gzip max_message_bytes: 1000000
Exemplo 2: Enviar logs com Collectd
O Collectd é um daemon que coleta periodicamente métricas de desempenho do sistema e de aplicações. Ele envia essas métricas para o SLS por meio do Write Kafka Plugin.
Instale o plugin Kafka e suas dependências. Por exemplo, no CentOS, execute sudo yum install collectd-write_kafka. Pacotes RPM estão disponíveis em Collectd-write_kafka.
-
Exemplo de configuração
Este exemplo utiliza saída em JSON. Outros formatos (Command, Graphite) estão documentados na documentação de configuração do Collectd.
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.
LoadPlugin write_kafka <Plugin write_kafka> Property "metadata.broker.list" "SLS_KAFKA_ENDPOINT" Property "security.protocol" "sasl_ssl" Property "sasl.mechanism" "PLAIN" Property "sasl.username" "SLS_PROJECT" Property "sasl.password" "SLS_PASSWORD" Property "broker.address.family" "v4" <Topic "SLS_LOGSTORE"> Format JSON Key "content" </Topic> </Plugin>
Exemplo 3: Enviar logs com Telegraf
O Telegraf é um agente leve baseado em Go para coleta, processamento e agregação de métricas de sistemas host, APIs de terceiros e serviços.
Modifique o arquivo de configuração do Telegraf para enviar logs ao SLS pelo protocolo Kafka.
-
Exemplo de configuração
-
Este exemplo utiliza saída em JSON. Outros formatos (Graphite, Carbon2) estão documentados em Formatos de saída do Telegraf.
NotaConfigure um caminho válido para tls_ca no Telegraf. Utilize o caminho do certificado raiz fornecido pelo servidor. No Linux, o caminho do certificado CA raiz geralmente é /etc/ssl/certs/ca-bundle.crt.
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.
# Kafka output plugin configuration [[outputs.kafka]] ## URLs of kafka brokers brokers = ["SLS_KAFKA_ENDPOINT"] ## Kafka topic for producer messages topic = "SLS_LOGSTORE" routing_key = "content" ## CompressionCodec represents the various compression codecs recognized by ## Kafka in messages. ## 0 : No compression ## 1 : Gzip compression ## 2 : Snappy compression ## 3 : LZ4 compression compression_codec = 1 ## Optional TLS Config tls_ca = "/etc/ssl/certs/ca-bundle.crt" tls_cert = "/etc/ssl/certs/ca-certificates.crt" # tls_key = "/etc/telegraf/key.pem" ## Use TLS but skip chain & host verification # insecure_skip_verify = false ## Optional SASL Config sasl_username = "SLS_PROJECT" sasl_password = "SLS_PASSWORD" ## Data format to output. ## https://github.com/influxdata/telegraf/blob/master/docs/DATA_FORMATS_OUTPUT.md data_format = "json" -
Exemplo 4: Enviar logs com Fluentd
O Fluentd é um coletor de logs open source.
Instale e configure o fluent-plugin-kafka para enviar logs ao SLS.
-
Exemplo de configuração
Este exemplo utiliza saída em JSON. Outros formatos estão documentados em Fluentd Formatter.
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.
<match **> @type kafka2 brokers SLS_KAFKA_ENDPOINT default_topic SLS_LOGSTORE default_message_key content sasl_over_ssl true use_event_time true username SLS_PROJECT password "SLS_PASSWORD" ssl_ca_certs_from_system true # ruby-kafka producer options max_send_retries 1000 required_acks 1 compression_codec gzip use_event_time true max_send_limit_bytes 2097152 <buffer hostlogs> flush_interval 10s </buffer> <format> @type json </format> </match>
Exemplo 5: Enviar logs com Logstash
O Logstash é um mecanismo de coleta de logs em tempo real e open source, capaz de coletar logs de diversas fontes.
O envio de logs pelo protocolo Kafka requer o Logstash 7.10.1 ou superior.
O Logstash possui um plugin de saída Kafka nativo. Como o SLS exige SASL_SSL, também é necessário configurar um certificado SSL e um arquivo JAAS.
-
Exemplo de configuração
-
Este exemplo utiliza saída em JSON. Outros formatos estão documentados em Logstash Codec.
NotaEste exemplo serve apenas para testes de conectividade. Remova a configuração de saída stdout em ambientes de produção.
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.
output { stdout { codec => rubydebug } kafka { topic_id => "SLS_LOGSTORE" bootstrap_servers => "SLS_KAFKA_ENDPOINT" security_protocol => "SASL_SSL" sasl_jaas_config => "org.apache.kafka.common.security.plain.PlainLoginModule required username='SLS_PROJECT' password='SLS_PASSWORD';" sasl_mechanism => "PLAIN" codec => "json" client_id => "kafka-logstash" } } -
Exemplo 6: Enviar logs com Fluent-bit
O Fluent-bit é um processador leve de logs e métricas que suporta o envio de logs para o SLS por meio do plugin de saída Kafka.
-
Exemplo de configuração
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.[Output] Name kafka Match * Brokers SLS_KAFKA_ENDPOINT Topics SLS_LOGSTORE Format json rdkafka.sasl.username SLS_PROJECT rdkafka.sasl.password SLS_PASSWORD rdkafka.security.protocol SASL_SSL rdkafka.sasl.mechanism PLAIN
Exemplo 7: Configurar o Vector para enviar logs pelo protocolo Kafka
O Vector é uma ferramenta leve e de alto desempenho para processamento de logs. Configure o Vector para gravar dados no SLS no modo compatível com Kafka conforme descrito abaixo.
-
Exemplo de configuração
Para os parâmetros com prefixo
SLS_neste exemplo, consulte Parâmetros de configuração.[sinks.aliyun_sls] type = "kafka" inputs = ["test_logs"] bootstrap_servers = "SLS_KAFKA_ENDPOINT" compression = "gzip" healthcheck = true topic = "SLS_LOGSTORE" encoding.codec = "json" sasl.enabled = true sasl.mechanism = "PLAIN" sasl.username = "SLS_PROJECT" sasl.password = "SLS_PASSWORD" tls.enabled = true
Exemplo 8: Usar um produtor Kafka para enviar logs
Java
-
Dependências
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>2.1.0</version> </dependency> -
Código de exemplo
package org.example; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class KafkaProduceExample { public static void main(String[] args) { // The configurations. Properties props = new Properties(); String project = "etl-shanghai-b"; String logstore = "testlog"; // Set this parameter to true if you want the content of the producer to be parsed as a JSON log. boolean parseJson = false; // An Alibaba Cloud account has permissions on all API operations, which poses high security risks. We recommend creating and using a RAM user to call API operations or perform routine O&M. To create a RAM user, log on to the RAM console. // The following example stores the AccessKey ID and AccessKey secret in environment variables. Alternatively, store them in a configuration file. // Do not hard-code the AccessKey ID and AccessKey secret in source code to avoid credential leaks. String accessKeyID = System.getenv("SLS_ACCESS_KEY_ID"); String accessKeySecret = System.getenv("SLS_ACCESS_KEY_SECRET"); String endpoint = "cn-shanghai.log.aliyuncs.com"; // Configure this parameter based on the endpoint of the project. String port = "10012"; // Use port 10012 for the Internet and port 10011 for the internal network. String hosts = project + "." + endpoint + ":" + port; String topic = logstore; if(parseJson) { topic = topic + ".json"; } props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, hosts); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put("security.protocol", "sasl_ssl"); props.put("sasl.mechanism", "PLAIN"); props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"" + project + "\" password=\"" + accessKeyID + "#" + accessKeySecret + "\";"); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); //Create a producer instance. KafkaProducer<String,String> producer = new KafkaProducer<>(props); //Send records. for(int i=0;i<1;i++){ String content = "{\"msg\": \"Hello World\"}"; // If needed, use the following method to set the timestamp of the message. // long timestamp = System.currentTimeMillis(); // ProducerRecord<String, String> record = new ProducerRecord<>(topic, null, timestamp, null, content); ProducerRecord<String, String> record = new ProducerRecord<>(topic, content); producer.send(record, (metadata, exception) -> { if (exception != null) { System.err.println("ERROR: Failed to send message: " + exception.getMessage()); exception.printStackTrace(); } else { System.out.println("Message sent successfully to topic: " + metadata.topic() + ", partition: " + metadata.partition() + ", offset: " + metadata.offset() + ", timestamp: " + metadata.timestamp()); } }); } producer.close(); } }
Python
-
Dependências
pip install confluent-kafka -
Código de exemplo
#!/bin/env python3 import time import os from confluent_kafka import Producer def delivery_report(err, msg): """ Called once for each message produced to indicate delivery result. Triggered by poll() or flush(). """ if err is not None: print('Message delivery failed: {}'.format(err)) else: print('Message delivered to {} [{}] at offset {}'.format(msg.topic(), msg.partition(), msg.offset())) def main(): project = "etl-shanghai-b" logstore = "testlog" parse_json = False # Get credentials from environment variables access_key_id = os.getenv("SLS_ACCESS_KEY_ID") access_key_secret = os.getenv("SLS_ACCESS_KEY_SECRET") endpoint = "cn-shanghai.log.aliyuncs.com" port = "10012" # Use port 10012 for the Internet and port 10011 for the internal network. hosts = f"{project}.{endpoint}:{port}" topic = logstore if parse_json: topic = topic + ".json" # Configure Kafka producer conf = { 'bootstrap.servers': hosts, 'security.protocol': 'sasl_ssl', 'sasl.mechanisms': 'PLAIN', 'sasl.username': project, 'sasl.password': f"{access_key_id}#{access_key_secret}", 'enable.idempotence': False, } # Create producer instance producer = Producer(conf) # Send message content = "{\"msg\": \"Hello World\"}" producer.produce(topic=topic, value=content.encode('utf-8'), #timestamp=int(time.time() * 1000), # (Optional) Set the timestamp of the record in milliseconds. callback=delivery_report) # Wait for any outstanding messages to be delivered and delivery report # callbacks to be triggered. producer.flush() if __name__ == '__main__': main()
Golang
-
Dependências
go get github.com/confluentinc/confluent-kafka-go/kafka -
Código de exemplo
package main import ( "fmt" "log" "os" // "time" "github.com/confluentinc/confluent-kafka-go/kafka" ) func main() { project := "etl-shanghai-b" logstore := "testlog" parseJson := false // Get credentials from environment variables accessKeyID := os.Getenv("SLS_ACCESS_KEY_ID") accessKeySecret := os.Getenv("SLS_ACCESS_KEY_SECRET") endpoint := "cn-shanghai.log.aliyuncs.com" port := "10012" // Use port 10012 for the Internet and port 10011 for the internal network. hosts := fmt.Sprintf("%s.%s:%s", project, endpoint, port) topic := logstore if parseJson { topic = topic + ".json" } // Configure Kafka producer config := &kafka.ConfigMap{ "bootstrap.servers": hosts, "security.protocol": "sasl_ssl", "sasl.mechanisms": "PLAIN", "sasl.username": project, "sasl.password": accessKeyID + "#" + accessKeySecret, "enable.idempotence": false, } // Create producer instance producer, err := kafka.NewProducer(config) if err != nil { log.Fatalf("Failed to create producer: %v", err) } defer producer.Close() // Send messages in batches. messages := []string{ "{\"msg\": \"Hello World 1\"}", "{\"msg\": \"Hello World 2\"}", "{\"msg\": \"Hello World 3\"}", } for _, content := range messages { err := producer.Produce(&kafka.Message{ TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny}, Value: []byte(content), //Timestamp: time.Now(), // Set the time if needed. }, nil) if err != nil { log.Printf("Failed to produce message: %v", err) } } // Enable a goroutine to listen for whether the producer successfully sends the message. go func() { for e := range producer.Events() { switch ev := e.(type) { case *kafka.Message: if ev.TopicPartition.Error != nil { fmt.Printf("Delivery failed: %v\n", ev.TopicPartition.Error) } else { fmt.Printf("Delivered message to topic %s [%d] at offset %v\n", *ev.TopicPartition.Topic, ev.TopicPartition.Partition, ev.TopicPartition.Offset) } } } }() producer.Flush(5 * 1000) }
Solução de problemas
A tabela a seguir lista erros comuns no envio de logs pelo protocolo Kafka. A referência completa está disponível na lista de erros.
|
Mensagem de erro |
Descrição |
Solução recomendada |
|
NetworkException |
Ocorreu um erro de rede. |
Aguarde um segundo e tente novamente. |
|
TopicAuthorizationException |
Falha na autenticação. |
O par de AccessKeys é inválido ou a conta não tem permissões de gravação no projeto ou Logstore especificado. Especifique um par de AccessKeys válido com as permissões necessárias. |
|
UnknownTopicOrPartitionException |
Este erro ocorre por um dos seguintes motivos:
|
Verifique se o projeto e o Logstore existem. Se o erro persistir, confirme se a região do projeto corresponde à região do endpoint. |
|
KafkaStorageException |
Ocorreu uma exceção no lado do servidor. |
Aguarde um segundo e tente novamente. |