Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Connect to ApsaraMQ for Kafka over SSL with PLAIN authentication

Última atualização: Sep 21, 2026

Para produzir e consumir mensagens em uma instância do ApsaraMQ for Kafka pela Internet, conecte-se ao endpoint ssl com autenticação SASL/PLAIN. O mecanismo PLAIN transmite credenciais em texto claro, portanto, combine-o sempre com a criptografia ssl (SASL_SSL) para proteger as credenciais em trânsito.

Todos os exemplos deste guia usam o SDK Java.

Espaços reservados

Reúna os seguintes valores antes de começar. Substitua esses espaços reservados em todas as configurações e exemplos de código.

Espaço reservado

Descrição

Onde encontrar

<bootstrap-servers>

Endpoint ssl da sua instância do ApsaraMQ for Kafka

Página Instance Details no console do ApsaraMQ for Kafka

<topic-name>

Nome do tópico

Página Topics no console do ApsaraMQ for Kafka

<group-id>

id do grupo de consumidores

Página Groups no console do ApsaraMQ for Kafka

<truststore-path>

Caminho absoluto para o certificado raiz ssl na sua máquina

Consulte a Etapa 2: baixe o certificado raiz ssl

<jaas-conf-path>

Caminho absoluto para o arquivo de configuração JAAS

Consulte a Etapa 3: crie o arquivo de configuração JAAS

<username>

Nome de usuário SASL

Página Instance Details (ACL desativada) ou suas credenciais de usuário SASL (ACL ativada)

<password>

Senha SASL

A mesma do nome de usuário

Pré-requisitos

Adicione dependências Java

Adicione as seguintes dependências ao seu arquivo pom.xml:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.6.0</version>
</dependency>
<dependency>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-log4j12</artifactId>
    <version>1.7.6</version>
</dependency>
Nota

Certifique-se de que a versão principal do kafka-clients corresponda à versão da sua instância do ApsaraMQ for Kafka. Localize a versão da instância na página Instance Details no console do ApsaraMQ for Kafka.

configure os arquivos de configuração

configure os seguintes arquivos antes de escrever o código do produtor e do consumidor.

Arquivo

Finalidade

log4j.properties

Configuração de saída de log

Certificado raiz ssl (.jks)

Âncora de confiança TLS para a conexão com o broker

kafka_client_jaas.conf

Credenciais SASL/PLAIN

kafka.properties

Endpoint do broker, tópico, grupo de consumidores e caminhos dos arquivos

JavaKafkaConfigurer.java

Classe auxiliar que carrega as propriedades e defina o caminho JAAS

Etapa 1: crie o arquivo de configuração do Log4j

crie um arquivo chamado log4j.properties:

# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements.  See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License.  You may obtain a copy of the License at
#
#    http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

log4j.rootLogger=INFO, STDOUT

log4j.appender.STDOUT=org.apache.log4j.ConsoleAppender
log4j.appender.STDOUT.layout=org.apache.log4j.PatternLayout
log4j.appender.STDOUT.layout.ConversionPattern=[%d] %p %m (%c)%n

Etapa 2: baixe o certificado raiz ssl

baixe o certificado raiz ssl e salve-o em um local acessível pela sua aplicação. Anote o caminho absoluto, pois você o usará como <truststore-path>.

Aviso

Em ambientes de produção, use o caminho absoluto completo para o arquivo truststore. Não o inclua dentro de um JAR.

Etapa 3: crie o arquivo de configuração JAAS

crie um arquivo chamado kafka_client_jaas.conf:

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="<username>"
  password="<password>";
};
Nota
  • Se o recurso de lista de controle de acesso (ACL) estiver desativado para a instância, localize as credenciais padrão do usuário Simple Authentication and Security Layer (SASL) na página Instance Details no console do ApsaraMQ for Kafka.

  • Se a ACL estiver ativada, certifique-se de que o usuário SASL seja do tipo PLAIN e tenha permissões para produzir e consumir mensagens. Para obter mais detalhes, consulte Grant permissions to SASL users.

Etapa 4: crie o arquivo de propriedades do Kafka

crie o arquivo kafka.properties:

## SSL endpoint (from the ApsaraMQ for Kafka console)
bootstrap.servers=<bootstrap-servers>
## Topic name (created in the ApsaraMQ for Kafka console)
topic=<topic-name>
## Consumer group ID (created in the ApsaraMQ for Kafka console)
group.id=<group-id>
## Absolute path to the SSL root certificate
ssl.truststore.location=<truststore-path>
## Absolute path to the JAAS configuration file
java.security.auth.login.config=<jaas-conf-path>

Etapa 5: crie o carregador de configuração

crie o arquivo JavaKafkaConfigurer.java para carregar o arquivo de propriedades e defina o caminho de configuração JAAS.

import java.util.Properties;

public class JavaKafkaConfigurer {

    private static Properties properties;

    public static void configureSasl() {
        // Skip if the JAAS config path is already set via -D flag or another method
        if (null == System.getProperty("java.security.auth.login.config")) {
            System.setProperty("java.security.auth.login.config",
                getKafkaProperties().getProperty("java.security.auth.login.config"));
        }
    }

    public synchronized static Properties getKafkaProperties() {
        if (null != properties) {
            return properties;
        }
        Properties kafkaProperties = new Properties();
        try {
            kafkaProperties.load(
                KafkaProducerDemo.class.getClassLoader().getResourceAsStream("kafka.properties"));
        } catch (Exception e) {
            e.printStackTrace();
        }
        properties = kafkaProperties;
        return kafkaProperties;
    }
}

Referência de propriedades ssl e SASL

As propriedades a seguir são comuns aos exemplos de produtor e consumidor. Elas configuram o transporte ssl e a autenticação SASL/PLAIN.

Propriedade

Valor

Descrição

security.protocol

SASL_SSL

Criptografa o tráfego com TLS e autentica com SASL

sasl.mechanism

PLAIN

Usa o mecanismo de autenticação PLAIN

ssl.truststore.location

<truststore-path>

Caminho para o certificado raiz ssl (arquivo .jks)

ssl.truststore.password

KafkaOnsClient

Senha padrão do truststore

ssl.endpoint.identification.algorithm

(string vazia)

Desativa a verificação de nome de host

Produza mensagens

crie o arquivo KafkaProducerDemo.java:

import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.Future;
import org.apache.kafka.clients.CommonClientConfigs;
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.clients.producer.RecordMetadata;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;

public class KafkaProducerDemo {

    public static void main(String args[]) {
        // Load JAAS configuration
        JavaKafkaConfigurer.configureSasl();

        // Load kafka.properties
        Properties kafkaProperties =  JavaKafkaConfigurer.getKafkaProperties();

        Properties props = new Properties();

        // --- SSL + SASL/PLAIN authentication ---
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
            kafkaProperties.getProperty("bootstrap.servers"));
        props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
            kafkaProperties.getProperty("ssl.truststore.location"));
        props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
        props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
        // Disable hostname verification
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");

        // --- Producer settings ---
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringSerializer");
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringSerializer");
        props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 30 * 1000);
        props.put(ProducerConfig.RETRIES_CONFIG, 5);
        props.put(ProducerConfig.RECONNECT_BACKOFF_MS_CONFIG, 3000);

        // Create a thread-safe producer (one per process is usually sufficient;
        // for higher throughput, create up to 5)
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(props);

        // Topic to send messages to
        String topic = kafkaProperties.getProperty("topic");
        String value = "this is the message's value";

        try {
            List<Future<RecordMetadata>> futures = new ArrayList<Future<RecordMetadata>>(128);
            for (int i = 0; i < 100; i++) {
                ProducerRecord<String, String> kafkaMessage =
                    new ProducerRecord<String, String>(topic, value + ": " + i);
                Future<RecordMetadata> metadataFuture = producer.send(kafkaMessage);
                futures.add(metadataFuture);
            }
            producer.flush();
            for (Future<RecordMetadata> future : futures) {
                try {
                    RecordMetadata recordMetadata = future.get();
                    System.out.println("Produce ok:" + recordMetadata.toString());
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        } catch (Exception e) {
            System.out.println("error occurred");
            e.printStackTrace();
        }
    }
}

Compile e execute o KafkaProducerDemo.java para enviar mensagens.

Consuma mensagens

Consumidor único

crie o arquivo KafkaConsumerDemo.java:

import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;

public class KafkaConsumerDemo {

    public static void main(String args[]) {
        // Load JAAS configuration
        JavaKafkaConfigurer.configureSasl();

        // Load kafka.properties
        Properties kafkaProperties =  JavaKafkaConfigurer.getKafkaProperties();

        Properties props = new Properties();

        // --- SSL + SASL/PLAIN authentication ---
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
            kafkaProperties.getProperty("bootstrap.servers"));
        props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
            kafkaProperties.getProperty("ssl.truststore.location"));
        props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
        props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
        // Disable hostname verification
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");

        // --- Consumer settings ---
        // Session timeout (default: 30s). If no heartbeat is received within this interval,
        // the broker removes the consumer from the group and triggers rebalancing.
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
        // Maximum bytes fetched per partition and per request.
        // Tune these values for Internet connections to control bandwidth.
        props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 32000);
        props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 32000);
        // Maximum number of records returned per poll.
        // Keep this low enough to process all records before the next poll;
        // otherwise the broker triggers rebalancing.
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 30);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.GROUP_ID_CONFIG,
            kafkaProperties.getProperty("group.id"));

        // Create a consumer instance
        KafkaConsumer<String, String> consumer =
            new org.apache.kafka.clients.consumer.KafkaConsumer<String, String>(props);

        // Subscribe to one or more topics
        List<String> subscribedTopics = new ArrayList<String>();
        subscribedTopics.add(kafkaProperties.getProperty("topic"));
        consumer.subscribe(subscribedTopics);

        // Poll loop
        while (true) {
            try {
                ConsumerRecords<String, String> records = consumer.poll(1000);
                // Process all records before the next poll.
                // For better throughput, offload processing to a separate thread pool.
                for (ConsumerRecord<String, String> record : records) {
                    System.out.println(
                        String.format("Consume partition:%d offset:%d",
                            record.partition(), record.offset()));
                }
            } catch (Exception e) {
                try {
                    Thread.sleep(1000);
                } catch (Throwable ignore) {
                }
                e.printStackTrace();
            }
        }
    }
}

Compile e execute o KafkaConsumerDemo.java para consumir mensagens.

Múltiplos consumidores

Para aumentar o throughput, execute várias threads de consumo no mesmo processo. O número total de consumidores em todos os processos não deve exceder o número de partições no tópico inscrito.

crie o arquivo KafkaMultiConsumerDemo.java:

import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
import org.apache.kafka.common.errors.WakeupException;

public class KafkaMultiConsumerDemo {

    public static void main(String args[]) throws InterruptedException {
        // Load JAAS configuration
        JavaKafkaConfigurer.configureSasl();

        // Load kafka.properties
        Properties kafkaProperties = JavaKafkaConfigurer.getKafkaProperties();

        Properties props = new Properties();

        // --- SSL + SASL/PLAIN authentication ---
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
            kafkaProperties.getProperty("bootstrap.servers"));
        props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
            kafkaProperties.getProperty("ssl.truststore.location"));
        props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
        props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
        // Disable hostname verification
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");

        // --- Consumer settings ---
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 30);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
            "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.GROUP_ID_CONFIG,
            kafkaProperties.getProperty("group.id"));

        // Start two consumer threads
        int consumerNum = 2;
        Thread[] consumerThreads = new Thread[consumerNum];
        for (int i = 0; i < consumerNum; i++) {
            KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(props);

            List<String> subscribedTopics = new ArrayList<String>();
            subscribedTopics.add(kafkaProperties.getProperty("topic"));
            consumer.subscribe(subscribedTopics);

            KafkaConsumerRunner kafkaConsumerRunner = new KafkaConsumerRunner(consumer);
            consumerThreads[i] = new Thread(kafkaConsumerRunner);
        }

        for (int i = 0; i < consumerNum; i++) {
            consumerThreads[i].start();
        }

        for (int i = 0; i < consumerNum; i++) {
            consumerThreads[i].join();
        }
    }

    static class KafkaConsumerRunner implements Runnable {
        private final AtomicBoolean closed = new AtomicBoolean(false);
        private final KafkaConsumer consumer;

        KafkaConsumerRunner(KafkaConsumer consumer) {
            this.consumer = consumer;
        }

        @Override
        public void run() {
            try {
                while (!closed.get()) {
                    try {
                        ConsumerRecords<String, String> records = consumer.poll(1000);
                        for (ConsumerRecord<String, String> record : records) {
                            System.out.println(
                                String.format("Thread:%s Consume partition:%d offset:%d",
                                    Thread.currentThread().getName(),
                                    record.partition(), record.offset()));
                        }
                    } catch (Exception e) {
                        try {
                            Thread.sleep(1000);
                        } catch (Throwable ignore) {
                        }
                        e.printStackTrace();
                    }
                }
            } catch (WakeupException e) {
                if (!closed.get()) {
                    throw e;
                }
            } finally {
                consumer.close();
            }
        }

        public void shutdown() {
            closed.set(true);
            consumer.wakeup();
        }
    }
}

Compile e execute o KafkaMultiConsumerDemo.java para consumir mensagens com várias threads.

Verifique o resultado

Depois de executar o produtor e o consumidor, verifique a saída do console.

Produtor -- uma saída bem-sucedida tem o seguinte formato:

Produce ok:send-and-subscribe-to-messages-by-using-an-ssl-endpoint-with-plain-authentication-0@0
Produce ok:send-and-subscribe-to-messages-by-using-an-ssl-endpoint-with-plain-authentication-0@1
...

Consumidor -- uma saída bem-sucedida tem o seguinte formato:

Consume partition:0 offset:0
Consume partition:0 offset:1
...

Se algum dos programas lançar uma exceção, consulte Troubleshoot ApsaraMQ for Kafka client errors.