Todos os produtos
Search
Central de documentação

DataHub:Compatibilidade com Kafka

Última atualização: Aug 26, 2026

O DataHub é totalmente compatível com o protocolo Apache Kafka. Use clientes Kafka nativos para ler dados do DataHub e gravar dados nele.

Mapeamento do Kafka para o DataHub

Tipos de tópico

O Kafka e o DataHub possuem mecanismos diferentes para dimensionamento de tópicos. Para garantir a compatibilidade com o comportamento do Kafka, defina o modo de dimensionamento como ONLY_EXTEND ao criar um tópico do DataHub. Nesse modo, é possível apenas adicionar novos shards a um tópico. Não há suporte para divisão, mesclagem ou remoção de shards.

Nomenclatura de tópicos

O nome de um tópico do Kafka mapeia para um projeto e um tópico do DataHub, separados por ponto (.). O mapeamento segue estas regras:

  • A parte antes do primeiro . corresponde ao projeto do DataHub; a parte posterior, ao tópico do DataHub. Por exemplo, test_project.test_topic mapeia para o projeto test_project e para o tópico test_topic.

  • Se o nome contiver vários caracteres ., apenas o primeiro . atua como separador. Os demais caracteres . e - são substituídos por _.

Partições

Cada shard ativo no DataHub corresponde a uma partição no Kafka. Por exemplo, se um tópico tiver cinco shards ativos, ele equivalerá a um tópico do Kafka com cinco partições. Ao gravar dados, especifique um ID de partição no intervalo [0, 4]. Caso nenhuma partição seja especificada, o cliente Kafka atribuirá uma automaticamente.

Tópico Tuple

Ao gravar dados do Kafka em um tópico Tuple, o esquema do tópico deve ter uma ou duas colunas, ambas do tipo STRING. Caso contrário, a operação de gravação falhará.

  • Se o esquema tiver uma coluna, apenas o valor será gravado e a chave será descartada.

  • Se o esquema tiver duas colunas, a primeira e a segunda corresponderão à chave e ao valor, respectivamente.

Além disso, não grave dados binários em um tópico Tuple, pois isso gera caracteres ilegíveis. Para armazenar dados binários, utilize um tópico Blob.

Tópico Blob

Ao gravar dados do Kafka em um tópico Blob, o valor da mensagem do Kafka é gravado no campo Blob. Se a chave da mensagem não for NULL, ela será gravada como um atributo do DataHub. O nome do atributo é __kafka_key__ e seu valor corresponde à chave da mensagem do Kafka.

Headers

Os headers do Kafka mapeiam para atributos do DataHub. Se o valor de um header for NULL, ele será ignorado e não gravado como atributo. Recomendamos não usar __kafka_key__ como chave de header para evitar conflitos com o nome de atributo interno dos tópicos Blob.

Grupos de consumidores

No DataHub, um ID de assinatura atua como um grupo de consumidores, mas só pode assinar um único tópico. Em contraste, um grupo de consumidores do Kafka pode assinar vários tópicos simultaneamente. Para oferecer compatibilidade com o modelo de assinatura do Kafka, o DataHub disponibiliza um recurso de grupo. Crie um grupo dentro de um projeto e vincule-o a vários tópicos para assinar todos eles sob um único grupo.

Um grupo gerencia internamente várias assinaturas do DataHub no servidor. Após vincular um tópico, o grupo cria automaticamente uma assinatura, que aparece na lista de assinaturas na página de detalhes do tópico. Não exclua essa assinatura manualmente. Essa ação impediria o grupo de assinar o tópico e causaria a perda de todos os offsets de consumo existentes.

Um único grupo pode assinar no máximo 50 tópicos. Para assinar mais tópicos, envie um ticket.

Parâmetros do Kafka

C=Consumidor, P=Produtor, S=Streams

Parâmetro

C/P/S

Valor

Obrigatório

Descrição

bootstrap.servers

*

Consulte a seção de endpoints do Kafka.

Sim

security.protocol

*

SASL_SSL

Sim

Para garantir a segurança dos dados, as conexões do Kafka para o DataHub usam criptografia SSL por padrão.

sasl.mechanism

*

PLAIN

Sim

Mecanismo de autenticação para credenciais AccessKey. Apenas PLAIN é suportado.

compression.type

P

LZ4

Não

Define o tipo de compactação das mensagens. Atualmente, apenas LZ4 é suportado.

group.id

C

project.topic:subId

ou

project.group

Sim

Se você usar o formato project.topic:subId, o ID deve corresponder ao tópico assinado. Caso contrário, os dados não poderão ser lidos. Recomendamos o uso do formato project.group.

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

Não

A estratégia padrão de atribuição de partições no Kafka é RangeAssignor. O DataHub atualmente suporta apenas esta estratégia. Não modifique este parâmetro.

session.timeout.ms

C/S

[60000, 180000]

Não

O padrão no Kafka é 10.000 ms. No entanto, como o DataHub exige um mínimo de 60.000 ms, esse valor é ajustado automaticamente para 60.000 ms.

heartbeat.interval.ms

C/S

Recomendado: 2/3 de session.timeout.ms

Não

O padrão do Kafka é 3.000 ms. Como osession.timeout.ms é ajustado para 60.000 ms, recomendamos definir explicitamente este valor como 40000 para evitar requisições frequentes de heartbeat.

application.id

S

project.topic:subId

ou

project.group

Sim

Se você usar o formato project.topic:subId, o ID deve corresponder ao tópico assinado. Caso contrário, os dados não poderão ser lidos. Recomendamos o uso do formato project.group.

Esta tabela lista os principais parâmetros a revisar ao usar um cliente Kafka com o DataHub. Outros parâmetros do lado do cliente, como retries e batch.size, comportam-se da mesma forma que no Kafka nativo. Parâmetros do lado do servidor não alteram o comportamento real do DataHub. Por exemplo, independentemente do valor de acks, o DataHub retorna uma confirmação somente após a gravação completa dos dados.

Endpoints do Kafka

Região

ID da região

Endpoint público

Endpoint ECS (rede clássica)

Endpoint ECS (VPC)

China (Hangzhou)

cn-hangzhou

dh-cn-hangzhou.aliyuncs.com:9092

dh-cn-hangzhou.aliyun-inc.com:9093

dh-cn-hangzhou-int-vpc.aliyuncs.com:9094

China (Shanghai)

cn-shanghai

dh-cn-shanghai.aliyuncs.com:9092

dh-cn-shanghai.aliyun-inc.com:9093

dh-cn-shanghai-int-vpc.aliyuncs.com:9094

China (Beijing)

cn-beijing

dh-cn-beijing.aliyuncs.com:9092

dh-cn-beijing.aliyun-inc.com:9093

dh-cn-beijing-int-vpc.aliyuncs.com:9094

China (Zhangjiakou)

cn-zhangjiakou

dh-cn-zhangjiakou.aliyuncs.com:9092

dh-cn-zhangjiakou.aliyun-inc.com:9093

dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094

China (Shenzhen)

cn-shenzhen

dh-cn-shenzhen.aliyuncs.com:9092

dh-cn-shenzhen.aliyun-inc.com:9093

dh-cn-shenzhen-int-vpc.aliyuncs.com:9094

Singapore

ap-southeast-1

dh-ap-southeast-1.aliyuncs.com:9092

dh-ap-southeast-1.aliyun-inc.com:9093

dh-ap-southeast-1-int-vpc.aliyuncs.com:9094

Malaysia (Kuala Lumpur)

ap-southeast-3

dh-ap-southeast-3.aliyuncs.com:9092

dh-ap-southeast-3.aliyun-inc.com:9093

dh-ap-southeast-3-int-vpc.aliyuncs.com:9094

Germany (Frankfurt)

eu-central-1

dh-eu-central-1.aliyuncs.com:9092

dh-eu-central-1.aliyun-inc.com:9093

dh-eu-central-1-int-vpc.aliyuncs.com:9094

China East 2 Finance

cn-shanghai-finance-1

dh-cn-shanghai-finance-1.aliyuncs.com:9092

dh-cn-shanghai-finance-1.aliyun-inc.com:9093

dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094

China (Hong Kong)

cn-hongkong

dh-cn-hongkong.aliyuncs.com:9092

dh-cn-hongkong.aliyun-inc.com:9093

dh-cn-hongkong-int-vpc.aliyuncs.com:9094

Criar um tópico

  1. Criar um tópico no console

    Ao criar um tópico, ative o Shard Expand Mode.

  2. Criar um tópico usando o SDK

    Não é possível criar tópicos pela API do Kafka. Use o SDK do DataHub e defina ExpandMode como ONLY_EXTEND. A versão necessária da dependência Maven é 2.19.0 ou posterior.

    Recomendamos configurar suas credenciais AccessKey ID e AccessKey Secret por meio de variáveis de ambiente para evitar codificá-las diretamente no código do projeto. O par de AccessKeys de uma conta Alibaba Cloud tem permissões para todas as operações de API. Para maior segurança, utilize um par de AccessKeys de um usuário RAM para acesso à API ou operações diárias, reduzindo assim o risco de vazamento de credenciais.

    datahub.endpoint=<yourEndpoint>
    datahub.accessId=<yourAccessKeyId>
    datahub.accessKey=<yourAccessKeySecret>
    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.19.0-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateTopic {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            int shardCount = 1;
            int lifeCycle = 7;
    
            try {
                datahubClient.createTopic("test_project", "test_topic", shardCount, lifeCycle, RecordType.BLOB, "comment", ExpandMode.ONLY_EXTEND);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Criar um grupo

  1. Criar um grupo no console

    Clique em Create Group e adicione os tópicos que deseja assinar na lista à direita. É possível modificar os tópicos vinculados após a criação do grupo. O grupo cria automaticamente uma assinatura, que aparece na página da lista de assinaturas do tópico.

  2. Criar um grupo usando o SDK

    A versão da dependência Maven deve ser 2.21.6-public ou posterior.

    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.21.6-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateGroup {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            List<String> topicList = new ArrayList<>();
            topicList.add("test_project.topic1");
            topicList.add("test_project.topic2");
            topicList.add("test_project.topic3");
    
            try {
                // Create a Kafka group.
                datahubClient.createKafkaGroup("test_project", "test_topic", "test comment");
    
                // Bind the topics to the group for subscription.
                datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Exemplo de produtor

O arquivo kafka_client_producer_jaas.conf

Crie um arquivo chamado kafka_client_producer_jaas.conf em qualquer diretório e adicione o seguinte conteúdo.

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="yourAccessKeyId"
  password="yourAccessKeySecret";
};

Dependência Maven

A versão do cliente Kafka deve ser 0.10.0.0 ou posterior. Recomendamos a versão 2.4.0.

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>

Código de exemplo

public class ProducerExample {
    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/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("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("compression.type", "lz4");

        String KafkaTopicName = "test_project.test_topic";
        Producer<String, String> producer = new KafkaProducer<String, String>(properties);

        try {
            List<Header> headers = new ArrayList<>();
            RecordHeader header1 = new RecordHeader("key1", "value1".getBytes());
            RecordHeader header2 = new RecordHeader("key2", "value2".getBytes());
            headers.add(header1);
            headers.add(header2);

            ProducerRecord<String, String> record = new ProducerRecord<>(KafkaTopicName, 0, "key", "Hello DataHub!", headers);

            // Sync send
            producer.send(record).get();

        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        } finally {
            producer.close();
        }
    }
}

Resultado

Após a execução bem-sucedida do código, faça uma amostragem dos dados para verificar o resultado.

Exemplo de consumidor

Para obter informações sobre como gerar o arquivo kafka_client_producer_jaas.conf e adicionar a dependência Maven, consulte o exemplo de produtor.

Quando um novo consumidor entra, a atribuição de shards leva de 10 a 20 segundos. Após a conclusão da atribuição, o consumidor pode começar a consumir dados.

Código de exemplo

Usando um grupo Kafka (recomendado)

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", "src/main/resources/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");
        // Set group.id to the project.group format.
        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");
        // By using a Kafka group, you can subscribe to multiple topics.
        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());
            }
        }
    }
}

Usando project.topic:subId

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.Collections;
import java.util.Properties;
public class ConsumerExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/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");
        // Set group.id to the project.topic:subId format.
        properties.put("group.id", "test_project.test_topic:1611039998153N71KM");
        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);

        // When using the project.topic:subId format, you can only subscribe to a single topic.
        kafkaConsumer.subscribe(Collections.singletonList("test_project.test_topic"));

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));

            for (ConsumerRecord<String, String> record : records) {
                System.out.println(record.toString());
            }
        }
    }
}

Resultado

Após a execução bem-sucedida do código, os dados consumidos serão exibidos no terminal.

ConsumerRecord(topic = test_project.test_topic, partition = 0, leaderEpoch = 0, offset = 0, LogAppendTime = 1611040892661, serialized key size = 3, serialized value size = 14, headers = RecordHeaders(headers = [RecordHeader(key = key1, value = [118, 97, 108, 117, 101, 49]), RecordHeader(key = key2, value = [118, 97, 108, 117, 101, 50])], isReadOnly = false), key = key, value = Hello DataHub!)

Neste exemplo, todos os registros de dados retornados em uma única requisição têm o mesmo LogAppendTime, que corresponde ao timestamp mais recente entre todos os registros desse lote.

Exemplo de Streams

Dependência Maven

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>2.4.0</version>
</dependency>

Código de exemplo

Este exemplo lê dados de um tópico de entrada em test_project, converte as strings de chave e valor para minúsculas e grava o resultado em um tópico de saída.

public class StreamExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/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.input:1611293595417QH0WL");
        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));
        }
    }
}

Resultado

Após iniciar a tarefa do Streams, a atribuição de shards leva cerca de um minuto. Em seguida, visualize o número de tarefas atuais 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.

currently assigned active tasks: [0_0, 0_1, 0_2]
currently assigned standby tasks: []
revoked active tasks: []
  revoked standby tasks: []

Após a atribuição dos shards, grave dados de teste como (AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF) no tópico de entrada. Em seguida, faça uma amostragem dos dados do tópico de saída para verificar se a gravação ocorreu corretamente.

Observações de uso

  • Não há suporte para transações e idempotência.

  • Clientes Kafka não podem criar tópicos automaticamente no DataHub. Crie o tópico antes de gravar dados nele.

  • Ao usar um ID de assinatura (project.topic:subid) como group.id, um consumidor só pode assinar um único tópico. Para assinar vários tópicos, utilize um grupo do DataHub.

  • O timestamp dos dados lidos por um consumidor é sempre o LogAppendTime, que indica quando os dados foram gravados no DataHub. Todos os registros em uma única requisição de busca compartilham o mesmo timestamp: o mais recente daquele lote. Isso significa que o timestamp de leitura pode ser posterior ao horário real de gravação.

  • Uma aplicação Streams suporta apenas um tópico de entrada, mas pode ter vários tópicos de saída.

  • Apenas tarefas Streams sem estado são suportadas.

  • As versões suportadas do Kafka variam de 0.10.0 a 2.4.0.

Perguntas frequentes

Desconexão ao gravar dados

Selector - [Producer clientId=producer-1] Connection with dh-cn-shenzhen.aliyuncs.com disconnected
java.io.EOFException
    at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:573)
    ...

As requisições de metadados do Kafka e as requisições de gravação de dados usam conexões diferentes.

O cliente estabelece primeiro uma conexão para buscar metadados. Em seguida, usa as informações do broker retornadas para estabelecer uma segunda conexão destinada à gravação de dados. Todas as requisições subsequentes são enviadas por essa segunda conexão.

A primeira conexão, agora ociosa, é fechada automaticamente pelo servidor após um tempo limite. Isso pode gerar um erro de desconexão nos logs. Ignore esse erro se os dados estiverem sendo gravados com sucesso.

Falha ao iniciar o cliente Kafka

Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed
Caused by: javax.net.ssl.SSLHandshakeException: No subject alternative names matching IP address 100.67.134.161 found

Adicione a seguinte propriedade à sua configuração: properties.put("ssl.endpoint.identification.algorithm", "");.

DisconnectException durante o consumo

[INFO][Consumer clientId=client-id, groupId=consumer-project.topic:subid] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: {}.
org.apache.kafka.common.errors.DisconnectException

O cliente Kafka precisa manter uma conexão TCP persistente com o servidor. Essa exceção geralmente é causada por instabilidade na rede. O cliente possui lógica de nova tentativa integrada, portanto, esse erro normalmente não afeta o consumo.