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 |
|
|
Endpoint ssl da sua instância do ApsaraMQ for Kafka |
Página Instance Details no console do ApsaraMQ for Kafka |
|
|
Nome do tópico |
Página Topics no console do ApsaraMQ for Kafka |
|
|
id do grupo de consumidores |
Página Groups no console do ApsaraMQ for Kafka |
|
|
Caminho absoluto para o certificado raiz ssl na sua máquina |
Consulte a Etapa 2: baixe o certificado raiz ssl |
|
|
Caminho absoluto para o arquivo de configuração JAAS |
Consulte a Etapa 3: crie o arquivo de configuração JAAS |
|
|
Nome de usuário SASL |
Página Instance Details (ACL desativada) ou suas credenciais de usuário SASL (ACL ativada) |
|
|
Senha SASL |
A mesma do nome de usuário |
Pré-requisitos
Uma instância, um tópico e um grupo de consumidores do ApsaraMQ for Kafka. Para obter mais detalhes, consulte Step 3: Create resources.
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>
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 |
|
|
Configuração de saída de log |
|
Certificado raiz ssl ( |
Âncora de confiança TLS para a conexão com o broker |
|
|
Credenciais SASL/PLAIN |
|
|
Endpoint do broker, tópico, grupo de consumidores e caminhos dos arquivos |
|
|
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>.
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>";
};
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 |
|
|
|
Criptografa o tráfego com TLS e autentica com SASL |
|
|
|
Usa o mecanismo de autenticação PLAIN |
|
|
|
Caminho para o certificado raiz ssl (arquivo |
|
|
|
Senha padrão do truststore |
|
|
(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.