Modo de compatibilidade do DataHub com Kafka
O DataHub agora é compatível com o protocolo Apache Kafka. Use um SDK padrão do Kafka para conectar-se ao DataHub e publicar ou assinar dados.
Mapeamento de conceitos entre DataHub e Kafka
|
Kafka |
DataHub |
|
topic |
Project.topic |
|
partition |
shard |
|
offset |
sequence |
Tópico do Kafka
O DataHub também possui tópicos, mas inclui uma camada adicional de recursos chamada projeto, inexistente no Kafka. Para garantir a compatibilidade, um tópico do Kafka corresponde à combinação de um projeto e um tópico do DataHub, unidos por um ponto (.). Por exemplo, se você tiver um projeto chamado test_project e um tópico chamado test_topic no DataHub, o nome do tópico Kafka correspondente será test_project.test_topic.
Partição do Kafka
Uma partição do Kafka equivale diretamente a um shard do DataHub. Ambos representam uma fila ordenada de dados.
Grupo de consumidores do Kafka
Um grupo de consumidores do DataHub funciona de maneira semelhante a um grupo de consumidores do Kafka. Crie um grupo de consumidores dentro de um projeto e associe-o aos tópicos que deseja assinar. O grupo pode assinar vários tópicos dentro desse projeto. Ao vincular um grupo de consumidores a um tópico, o sistema cria automaticamente uma assinatura e a exibe na página de lista de assinaturas do tópico. Excluir essa assinatura impede que o grupo assine o tópico e apaga todos os offsets de consumo anteriores.
Como o grupo de consumidores é um sub-recurso de um projeto, especifique-o junto com o projeto. Por exemplo, se o seu projeto do DataHub for test_project e o seu grupo de consumidores for test_group, o nome do grupo Kafka será test_project.test_group.
Na caixa de diálogo Create consumer group, insira um Name e uma Description, selecione os tópicos para assinar na área Topic, clique no botão > para adicioná-los à lista selecionada à direita e clique em Create.
Cada grupo pode assinar no máximo 50 tópicos. Caso precise assinar mais, envie um ticket.
Registro do Kafka
Um registro do Kafka utiliza o formato chave-valor. Os registros do DataHub estão disponíveis em dois formatos: Tuple (dados fortemente estruturados) e Blob (dados binários). Um registro do Kafka consiste em duas partes: um cabeçalho e um par chave-valor.
Cabeçalho do Kafka
Um cabeçalho do Kafka adiciona metadados aos registros, equivalente a um atributo do DataHub. Ao usar um cliente Kafka para gravar dados com cabeçalhos, as informações do cabeçalho são armazenadas como atributos do DataHub. Se o valor de um cabeçalho Kafka for NULL, o sistema ignorará o cabeçalho correspondente. Não use "__kafka_key__" como chave de cabeçalho, pois é uma chave interna reservada.
Dados chave-valor do Kafka
Se um tópico do DataHub for do tipo Tuple, sua chave corresponderá à primeira coluna String e seu valor corresponderá à segunda coluna String.
Se um tópico do DataHub for do tipo Blob, seu valor corresponderá ao conteúdo dos dados e sua chave será colocada em um atributo com o formato <"__kafka_key__", key>.
Offset do Kafka
Um offset do Kafka é um inteiro de 64 bits que começa em 0 e se autoincrementa dentro de cada partição, fornecendo uma posição única para cada registro. Ele serve principalmente para rastrear o progresso do consumidor. Uma sequência do DataHub equivale diretamente a um offset do Kafka.
Limites
O DataHub não oferece suporte a transações Kafka, idempotência, SchemaRegistry ou Log Compaction.
Início rápido
Use a versão recomendada do cliente Kafka 2.4.0. Ela é compatível com as versões 0.10.0 a 4.0.
Crie os recursos correspondentes de Project, Topic e Group no DataHub.
Altere o método de autenticação de segurança para SASL/SSL, use o mecanismo PLAIN para SASL e configure seu par de AccessKey da Alibaba Cloud (AK/SK).
Altere as informações do broker Kafka para o endpoint do DataHub (consulte a lista de domínios de service).
Criação de recursos
Primeiro, faça login no console do DataHub e crie um projeto.
Na caixa de diálogo New Project, defina Name como test_kafka_project e Description como test kafka.
Criar um tópico
Na página New Topic do console do DataHub, selecione Direct Create como método de criação, insira test_kafka como nome e selecione TUPLE como tipo. Na seção Schema details, adicione dois campos chamados key e value, defina o tipo de ambos como STRING e selecione Allow NULL. Defina Shard Count como 1 e Lifecycle como 3, ative Shard Scaling Mode e desative Multi-version. Insira uma descrição e clique em Create.
Crie um grupo e vincule os tópicos que precisa consumir. Modifique a lista de tópicos vinculados após a criação do grupo. Se não precisar consumir mensagens, pule esta etapa.
No painel New Group, insira um Name e uma Description, use a lista de transferência para mover os tópicos a serem consumidos da lista disponível à esquerda para a lista selecionada à direita e clique em Create.
Método de autenticação
Crie um arquivo chamado kafka_client_producer_jaas.conf e salve-o em qualquer caminho com o seguinte conteúdo:
KafkaClient {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="accessId"
password="accessKey";
};
Crie um arquivo chamado kafka_client_producer_jaas.conf e salve-o em qualquer diretório com o seguinte conteúdo:
KafkaClient {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="YOUR_ACCESS_ID"
password="YOUR_ACCESS_KEY";
};
Exemplo de produtor
Arquivo pom de clientes kafka
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>
Arquivo pom.xml do cliente Kafka:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>
ProducerExample
Exemplo de consumidor
ConsumerExample
package com.aliyun.datahub.kafka.demo;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
public class ConsumerExample2 {
static {
System.setProperty("java.security.auth.login.config", "/path/xxx/kafka_client_producer_jaas.conf");
}
public static void main(String[] args) {
Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("group.id", "test_project.test_kafka_group");
properties.put("auto.offset.reset", "earliest");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("ssl.endpoint.identification.algorithm", "");
properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);
List<String> topicList = new ArrayList<>();
topicList.add("test_project.test_topic1");
topicList.add("test_project.test_topic2");
topicList.add("test_project.test_topic3");
kafkaConsumer.subscribe(topicList);
while (true) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.toString());
}
}
}
}
ConsumerExample
package com.aliyun.datahub.kafka.demo;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
public class ConsumerExample {
static {
System.setProperty("java.security.auth.login.config", "/path/to/your/kafka_client_producer_jaas.conf");
}
public static void main(String[] args) {
Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("group.id", "test_project.test_kafka_group");
properties.put("auto.offset.reset", "earliest");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("ssl.endpoint.identification.algorithm", "");
properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);
List<String> topicList = new ArrayList<>();
topicList.add("test_project.test_topic1");
topicList.add("test_project.test_topic2");
topicList.add("test_project.test_topic3");
kafkaConsumer.subscribe(topicList);
while (true) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.toString());
}
}
}
}
Exemplos de Streams
Este código lê dados de input em test_project, converte as strings de chave e valor para minúsculas e regrava os dados em output.
public class StreamExample {
static {
System.setProperty("java.security.auth.login.config", "/path/xxx/kafka_client_producer_jaas.conf");
}
public static void main(final String[] args) {
final String input = "test_project.input";
final String output = "test_project.output";
final Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("application.id", "test_project.test_kafka_group");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("auto.offset.reset", "earliest");
final StreamsBuilder builder = new StreamsBuilder();
TestMapper testMapper = new TestMapper();
builder.stream(input, Consumed.with(Serdes.String(), Serdes.String()))
.map(testMapper)
.to(output, Produced.with(Serdes.String(), Serdes.String()));
final KafkaStreams streams = new KafkaStreams(builder.build(), properties);
final CountDownLatch latch = new CountDownLatch(1);
Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
@Override
public void run() {
streams.close();
latch.countDown();
}
});
try {
streams.start();
latch.await();
} catch (final Throwable e) {
System.exit(1);
}
System.exit(0);
}
static class TestMapper implements KeyValueMapper<String, String, KeyValue<String, String>> {
@Override
public KeyValue<String, String> apply(String s, String s2) {
return new KeyValue<>(StringUtils.lowerCase(s), StringUtils.lowerCase(s2));
}
}
}
Este exemplo lê dados do tópico de entrada no projeto test_project, converte as strings de chave e valor para minúsculas e grava o resultado no tópico de saída.
public class StreamExample {
static {
System.setProperty("java.security.auth.login.config", "/path/to/your/kafka_client_producer_jaas.conf");
}
public static void main(final String[] args) {
final String input = "test_project.input";
final String output = "test_project.output";
final Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("application.id", "test_project.test_kafka_group");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("auto.offset.reset", "earliest");
final StreamsBuilder builder = new StreamsBuilder();
TestMapper testMapper = new TestMapper();
builder.stream(input, Consumed.with(Serdes.String(), Serdes.String()))
.map(testMapper)
.to(output, Produced.with(Serdes.String(), Serdes.String()));
final KafkaStreams streams = new KafkaStreams(builder.build(), properties);
final CountDownLatch latch = new CountDownLatch(1);
Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
@Override
public void run() {
streams.close();
latch.countDown();
}
});
try {
streams.start();
latch.await();
} catch (final Throwable e) {
System.exit(1);
}
System.exit(0);
}
static class TestMapper implements KeyValueMapper<String, String, KeyValue<String, String>> {
@Override
public KeyValue<String, String> apply(String s, String s2) {
return new KeyValue<>(StringUtils.lowerCase(s), StringUtils.lowerCase(s2));
}
}
}
Após iniciar o job do Streams, a atribuição de shards leva cerca de um minuto. Em seguida, visualize a contagem de tarefas no console. O número de tarefas corresponde ao número de shards no tópico de entrada. Neste exemplo, o tópico de entrada possui três shards.
O resultado da amostragem de saída mostra que os dados foram gravados corretamente. Após o processamento pelo TestMapper na tarefa do Streams, a entrada em maiúsculas foi convertida para saída em minúsculas: as chaves são aaaa, cccc e eeee, e os valores correspondentes são bbbb, dddd e ffff.
Migrar dados de Kafka autogerenciado para o DataHub
Troque o endereço do broker. Para mais informações, consulte a lista de nomes de domínio abaixo.
Crie recursos no DataHub e modifique os nomes dos recursos no seu código. Para mais detalhes, veja a seção de criação de recursos acima.
Mude o método de autenticação para SASL_SSL. O mecanismo de autenticação é PLAIN. Defina username como seu AccessKey ID da Alibaba Cloud (AK) e password como seu AccessKey Secret (SK).
Apêndice
Visão geral da configuração
C=Consumidor, P=Produtor, S=Streams
|
Parâmetro |
C/P/S |
Valor recomendado |
Obrigatório |
Descrição |
|
bootstrap.servers |
* |
Consulte a lista de nomes de domínio. |
Sim |
|
|
security.protocol |
* |
SASL_SSL |
Sim |
Para garantir a transmissão segura de dados, a gravação de dados do Kafka para o DataHub usa criptografia SSL por padrão. |
|
sasl.mechanism |
* |
PLAIN |
Sim |
Autenticação via AccessKey. Apenas PLAIN é suportado. |
|
compression.type |
P |
LZ4 |
Não |
Especifica se deve ativar a compactação para transmissão de dados. Atualmente, apenas LZ4 é suportado. |
|
enable.idempotence |
P |
false |
Não |
Clientes Kafka versão 3.0.1 e posteriores ativam idempotência por padrão. Como o DataHub não suporta idempotência, desative manualmente esse recurso. Esta configuração não é necessária para clientes anteriores à versão 3.0.1. |
|
group.id |
C |
project.group |
Sim |
|
|
partition.assignment.strategy |
C |
org.apache.kafka.clients.consumer.RangeAssignor |
Não |
O valor padrão para Kafka é RangeAssignor. O DataHub atualmente suporta apenas RangeAssignor. Não altere esta configuração. |
|
session.timeout.ms |
C/S |
[60000, 180000] |
Não |
O valor padrão no Kafka é 10000. No entanto, o DataHub impõe um mínimo de 60000, então o valor assume 60000 por padrão. |
|
heartbeat.interval.ms |
C/S |
Dois terços do valor de session.timeout.ms. |
Não |
O valor padrão no Kafka é 3000. Como |
|
application.id |
S |
project.topic:subId ou project.group |
Sim |
Se usar project.topic:subId, ele deve corresponder ao tópico assinado. Caso contrário, os dados não poderão ser lidos. Recomendamos usar project.group. |
Lista de nomes de domínio de service
|
Nome da região |
Região |
Endpoint público |
Endpoint VPC ECS |
|
China East 1 (Hangzhou) |
cn-hangzhou |
dh-cn-hangzhou.aliyuncs.com:9092 |
dh-cn-hangzhou-int-vpc.aliyuncs.com:9094 |
|
China East 2 (Shanghai) |
cn-shanghai |
dh-cn-shanghai.aliyuncs.com:9092 |
dh-cn-shanghai-int-vpc.aliyuncs.com:9094 |
|
China North 2 (Beijing) |
cn-beijing |
dh-cn-beijing.aliyuncs.com:9092 |
dh-cn-beijing-int-vpc.aliyuncs.com:9094 |
|
China (Ulanqab) |
cn-wulanchabu |
dh-cn-wulanchabu.aliyuncs.com:9092 |
dh-cn-wulanchabu-int-vpc.aliyuncs.com:9094 |
|
China South 1 (Shenzhen) |
cn-shenzhen |
dh-cn-shenzhen.aliyuncs.com:9092 |
dh-cn-shenzhen-int-vpc.aliyuncs.com:9094 |
|
China North 3 (Zhangjiakou) |
cn-zhangjiakou |
dh-cn-zhangjiakou.aliyuncs.com:9092 |
dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094 |
|
Asia Pacific SE 1 (Singapore) |
ap-southeast-1 |
dh-ap-southeast-1.aliyuncs.com:9092 |
dh-ap-southeast-1-int-vpc.aliyuncs.com:9094 |
|
Asia Pacific SE 3 (Kuala Lumpur) |
ap-southeast-3 |
dh-ap-southeast-3.aliyuncs.com:9092 |
dh-ap-southeast-3-int-vpc.aliyuncs.com:9094 |
|
Europe Central 1 (Frankfurt) |
eu-central-1 |
dh-eu-central-1.aliyuncs.com:9092 |
dh-eu-central-1-int-vpc.aliyuncs.com:9094 |
|
China East 2 (Shanghai) Finance |
cn-shanghai-finance-1 |
dh-cn-shanghai-finance-1.aliyuncs.com:9092 |
dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094 |
|
China (Hong Kong) |
cn-hongkong |
dh-cn-hongkong.aliyuncs.com:9092 |
dh-cn-hongkong-int-vpc.aliyuncs.com:9094 |
Compatibilidade da API Kafka
A documentação oficial do Kafka lista todas as APIs disponíveis. A tabela a seguir apresenta as APIs que o DataHub suporta no modo de compatibilidade com Kafka.
|
API |
Descrição |
|
Produce |
Grava dados. |
|
Fetch |
Lê dados. |
|
ListOffsets |
Obtém offsets com base no tempo. Versões anteriores retornam uma lista de offsets, enquanto versões posteriores retornam um único offset. |
|
Metadata |
Obtém metadados para operações de leitura e gravação. |
|
OffsetCommit |
Confirma um offset. |
|
OffsetFetch |
Obtém um offset. |
|
FindCoordinator |
Localiza o broker que hospeda o coordenador e retorna seu endereço IP virtual (VIP). |
|
JoinGroup |
Entra em um grupo. |
|
Heartbeat |
Envia um heartbeat. |
|
LeaveGroup |
Sai de um grupo. |
|
SyncGroup |
Usado pelo líder do grupo para enviar o plano de atribuição de partições e por todos os membros para recebê-lo. |
|
SaslHandshake |
Gerencia o handshake do Simple Authentication and Security Layer (SASL) para autenticação. |
|
ApiVersions |
Obtém todas as APIs disponíveis. |
|
CreateTopics |
Cria um tópico. |
|
DeleteTopics |
Exclui um tópico. |
|
OffsetForLeaderEpoch |
Obtém o offset mais recente. |
|
SaslAuthenticate |
Autentica o cliente usando SASL. |
|
CreatePartitions |
Adiciona uma partição. |
|
DeleteGroups |
Exclui um grupo. |
|
OffsetDelete |
Exclui um offset. |