Todos os produtos
Search
Central de documentação

DataHub:Modo de compatibilidade do DataHub com Kafka

Última atualização: Jul 10, 2026

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

  1. Use a versão recomendada do cliente Kafka 2.4.0. Ela é compatível com as versões 0.10.0 a 4.0.

  2. Crie os recursos correspondentes de Project, Topic e Group no DataHub.

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

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

  1. Troque o endereço do broker. Para mais informações, consulte a lista de nomes de domínio abaixo.

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

  3. 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 session.timeout.ms tem como padrão 60000, defina explicitamente este valor como 40000 para evitar solicitações frequentes de heartbeat.

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.