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.
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
-
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
-
Acesse o diretório
configpara modificar os arquivos de configuração de origem e destino. Os arquivos são descritos da seguinte forma: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.StringDeserializerb.
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=falseDescrição dos parâmetros
Para obter uma lista de nomes de domínio para o parâmetro
bootstrap.serversde 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.configespecifica o módulo de login e as credenciais necessários para autenticação SASL. SubstituaAccessKey IDeAccessKey Secretpelas informações do seu AccessKey.Para mais detalhes sobre os itens de configuração, veja Compatibilidade com Kafka.
-
Configure o arquivo
topic-map.propertiespara 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
-
Defina as configurações de log no arquivo
log4j.properties.Crie um arquivo
log4j.properties.-
Utilize o modelo abaixo para sua configuração:
-
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
-
-
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 & -
Analise o log em busca de erros. A saída abaixo indica que a inicialização foi bem-sucedida:
-
[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
-
-
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á comoACTIVEe 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.