Todos os produtos
Search
Central de documentação

DataHub:Migrar dados para o DataHub com o Kafka Mirror Maker

Última atualização: Jun 28, 2026

Utilize a ferramenta Kafka Mirror Maker para migrar dados do Kafka para o DataHub.

Pré-requisitos

Verifique se você já criou um Project e um Topic. Para mais informações, consulte Criar um Topic.

Nota
  • A migração de dados é suportada apenas do Kafka para o DataHub. Não há suporte para migrar dados do DataHub para o Kafka.

  • O DataHub não oferece suporte a transações ou idempotência. Desative a idempotência na configuração do DataHub de destino.

Procedimento

  1. Faça upload do pacote kafka_mirror_datahub.tgz para o seu servidor Kafka de origem e extraia-o.

    tar -zxvf kafka_mirror_datahub.tgz
  1. Acesse o diretório config para modificar os arquivos de configuração de origem e destino. Os arquivos são descritos da seguinte forma:

    1. consumer.properties: arquivo de configuração do cluster Kafka de origem.

    # The server configuration for the source Kafka cluster
    bootstrap.servers=xx:9092
    # The consumer group ID
    group.id=test-consumer-group
    auto.offset.reset=earliest
    session.timeout.ms=60000
    heartbeat.interval.ms=40000
    ssl.endpoint.identification.algorithm=
    key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
    value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

    b. producer.properties: arquivo de configuração do serviço DataHub de destino.

    bootstrap.servers=dh-cn-zhangjiakou-pre.aliyuncs.com:9092
    sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required\nusername=\"AccessKey ID\"\npassword=\"AccessKey Secret\";
    security.protocol=SASL_SSL
    sasl.mechanism=PLAIN
    key.serializer=org.apache.kafka.common.serialization.StringSerializer
    value.serializer=org.apache.kafka.common.serialization.StringSerializer
    compression.type=lz4
    # Disable idempotence
    enable.idempotence=false

    Descrição dos parâmetros

    • Para obter uma lista de nomes de domínio para o parâmetro bootstrap.servers de destino, consulte Compatibilidade com Kafka. O endpoint neste exemplo, dh-cn-zhangjiakou-pre.aliyuncs.com:9092, corresponde à região China (Zhangjiakou).

    • O parâmetro sasl.jaas.config especifica o módulo de login e as credenciais necessários para autenticação SASL. Substitua AccessKey ID e AccessKey Secret pelas informações do seu AccessKey.

    • Para mais detalhes sobre os itens de configuração, veja Compatibilidade com Kafka.

  1. Configure o arquivo topic-map.properties para mapeamento de tópicos.

    Este arquivo mapeia tópicos do Kafka de origem para Topics do DataHub de destino. Cada linha representa uma regra de mapeamento, onde o lado esquerdo indica o nome do tópico de origem e o direito, o destino. Um Topic do DataHub de destino deve ser especificado no formato Project.Topic. Insira cada regra de mapeamento em uma nova linha.

    topicname=testproject.testtopic
    topicname1=testproject1.testtopic1
  1. Defina as configurações de log no arquivo log4j.properties.

    1. Crie um arquivo log4j.properties.

    2. Utilize o modelo abaixo para sua configuração:

      1. log4j.rootLogger=INFO, stdout, file
        log4j.appender.stdout=org.apache.log4j.ConsoleAppender
        log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
        log4j.appender.stdout.layout.ConversionPattern=[%d{yyyy-MM-dd HH:mm:ss}] [%p] %m (%c:%L)%n
        log4j.appender.file=org.apache.log4j.DailyRollingFileAppender
        log4j.appender.file.DatePattern='.'yyyy-MM-dd
        log4j.appender.file.File=/opt/logs/mm1.log
        log4j.appender.file.layout=org.apache.log4j.PatternLayout
        log4j.appender.file.layout.ConversionPattern=[%d{yyyy-MM-dd HH:mm:ss}] [%p] %m (%c:%L)%n
        log4j.logger.kafka=INFO
        log4j.logger.org.apache.kafka=INFO
        log4j.logger.kafka.tools.MirrorMaker=INFO
        log4j.logger.org.apache.zookeeper=WARN
  2. Execute o script de migração.

    Execute o script a seguir no diretório raiz da instalação do Kafka e verifique a saída de log.

    Descrição dos parâmetros

    • --consumer.config: arquivo de configuração do cluster Kafka de origem.

    • --producer.config: arquivo de configuração do serviço DataHub de destino.

    • --whitelist: nomes dos tópicos de origem. Para especificar vários tópicos, separe-os com uma barra vertical (|), por exemplo, topicA|topicB|topicC.

    • --topic.mapping.file: arquivo de configuração para mapeamento de tópicos.

    • KAFKA_LOG4J_OPTS: caminho para o arquivo de configuração de log.

    nohup KAFKA_LOG4J_OPTS="log4j.properties"  bin/kafka-mirror-maker.sh   --consumer.config config/consumer.properties   --producer.config config/producer.properties   --whitelist "mirrortest"   --topic.mapping.file /opt/kafka_2.12-3.7.2/config/topic-map.properties ... > /dev/null 2>&1 &
  3. Analise o log em busca de erros. A saída abaixo indica que a inicialização foi bem-sucedida:

    1. [2025-08-06 17:27:41] [INFO] Registered kafka:type=kafka.Log4jController MBean (kafka.utils.Log4jControllerRegistration$:31)
      [2025-08-06 17:27:41] [INFO] Starting mirror maker (kafka.tools.MirrorMaker$:62)
      [2025-08-06 17:27:41] [INFO] Loaded topic mappings: mirrortest -> test_suyang.mirror (kafka.tools.MirrorMaker$:62)
      [2025-08-06 17:27:41] [INFO] ProducerConfig values:
      	acks = -1
      	batch.size = 16384
      	bootstrap.servers = [xxx]
      	buffer.memory = 33554432
      	client.dns.lookup = use_all_dns_ips
      	client.id = producer-1
      	compression.type = lz4
      	connections.max.idle.ms = 540000
      	delivery.timeout.ms = 2147483647
      	enable.idempotence = false
      	interceptor.classes = []
      	internal.auto.downgrade.txn.commit = false
      	key.serializer = class org.apache.kafka.common.serialization.ByteArraySerializer
      	linger.ms = 0
      	max.block.ms = 9223372036854775807
      	max.in.flight.requests.per.connection = 1
      	max.request.size = 1048576
      	metadata.max.age.ms = 300000
      	metadata.max.idle.ms = 300000
  4. Faça login no console do DataHub para confirmar que os dados foram gravados corretamente.

    Na página Data Bus, localize o Topic de destino dentro do Project correspondente, como kafkatest/test. Na aba Shard List, verifique se o status do shard está como ACTIVE e se há um timestamp de dados recente. Clique em Sample na coluna Actions. No painel exibido, selecione um Shard ID, especifique a quantidade de registros a serem recuperados e clique em Sample. Na tabela de pré-visualização de dados, confirme que a gravação no DataHub ocorreu com sucesso.