Todos os produtos
Search
Central de documentação

Simple Log Service:Envio de logs pelo protocolo Kafka

Última atualização: Jul 03, 2026

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:

  • AliyunLogFullAccess

    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

    1. 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.

      Nota

      No script, substitua project_name pelo 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"
              }
          ]
      }
    2. 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 Nome do projeto.Endpoint:Porta. Configure este parâmetro com base no endpoint do projeto. Para mais informações, consulte Endpoint.

  • Rede interna: porta 10011. Exemplo: Nome do projeto.cn-hangzhou-intranet.log.aliyuncs.com:10011.

  • Internet: porta 10012. Exemplo: Nome do projeto.cn-hangzhou.log.aliyuncs.com:10012.

aliyun-project-test é o nome do projeto. cn-hangzhou-xxx.aliyuncs.com é o endpoint. 10011 e 10012 são as portas para rede interna e Internet, respectivamente.

  • Rede interna: aliyun-project-test.cn-hangzhou-intranet.log.aliyuncs.com:10011.

  • Internet: aliyun-project-test.cn-hangzhou.log.aliyuncs.com:10012.

SLS_PROJECT

Nome do projeto

Nome do projeto SLS.

aliyun-project-test

SLS_LOGSTORE

Nome do Logstore

Nome do Logstore. Se você adicionar .json ao final do nome, o SLS tentará analisar o conteúdo do log como JSON.

Por exemplo, se o nome do Logstore for test-logstore:

  • Se o valor for test-logstore, o conteúdo do log enviado será armazenado no campo content.

  • Se o valor for test-logstore.json, o conteúdo enviado será analisado como um log JSON. As chaves do primeiro nível dos dados JSON serão usadas como nomes de campo e os valores correspondentes como valores de campo.

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 #.

  • AccessKey ID: o AccessKey ID de uma conta Alibaba Cloud ou de um usuário RAM.

  • AccessKey secret: o AccessKey secret de uma conta Alibaba Cloud ou de um usuário RAM.

LTAI#yourAccessKeySecret

Nota

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

    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.

      Nota

      Configure 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

    <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.

Nota

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.

      Nota

      Este 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:

  • O projeto ou Logstore especificado não existe.

  • A região do projeto difere da região especificada no endpoint.

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.