Anda dapat menggunakan alat Kafka Mirror Maker untuk memigrasikan data dari Kafka ke DataHub.
Prasyarat
Pastikan Anda telah membuat Proyek dan Topik. Untuk informasi selengkapnya, lihat Buat Topik.
-
Anda hanya dapat memigrasikan data dari Kafka ke DataHub. Migrasi data dari DataHub ke Kafka tidak didukung.
-
DataHub tidak mendukung transaksi atau idempotence. Anda harus menonaktifkan idempotence dalam konfigurasi untuk DataHub tujuan.
Prosedur
-
Unggah paket kafka_mirror_datahub.tgz ke server Kafka sumber Anda dan ekstrak paket tersebut.
tar -zxvf kafka_mirror_datahub.tgz
-
Buka direktori
configuntuk mengubah file konfigurasi sumber dan tujuan. File-file tersebut dijelaskan sebagai berikut:-
consumer.properties: File konfigurasi untuk kluster Kafka sumber.
# Konfigurasi server untuk kluster Kafka sumber bootstrap.servers=xx:9092 # ID kelompok konsumen 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: File konfigurasi untuk layanan DataHub tujuan.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 # Nonaktifkan idempotence enable.idempotence=falsePenjelasan parameter
-
Untuk daftar nama domain
bootstrap.serverstujuan, lihat Kompatibilitas Kafka. Titik akhir pada contoh ini,dh-cn-zhangjiakou-pre.aliyuncs.com:9092, sesuai dengan wilayah China (Zhangjiakou). -
Parameter
sasl.jaas.configmenentukan modul login dan kredensial yang diperlukan untuk otentikasi SASL. GantiAccessKey IDdanAccessKey Secretdengan Informasi AccessKey Anda. -
Untuk informasi lebih lanjut tentang item konfigurasi, lihat Kompatibilitas Kafka.
-
-
Konfigurasikan file
topic-map.propertiesuntuk pemetaan topik.File ini memetakan topik Kafka sumber ke Topik DataHub tujuan. Setiap baris merepresentasikan aturan pemetaan di mana sisi kiri adalah nama topik sumber dan sisi kanan adalah tujuan. Topik DataHub tujuan ditentukan dalam format
Proyek.Topik. Tempatkan setiap aturan pemetaan pada baris baru.topicname=testproject.testtopic topicname1=testproject1.testtopic1
-
Konfigurasikan file
log4j.propertiesuntuk logging.-
Buat file
log4j.properties. -
Gunakan templat berikut untuk konfigurasi Anda:
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
-
-
Jalankan skrip migrasi.
Jalankan skrip berikut dari direktori root instalasi Kafka Anda dan periksa output log.
Penjelasan parameter
-
--consumer.config: File konfigurasi untuk kluster Kafka sumber. -
--producer.config: File konfigurasi untuk layanan DataHub tujuan. -
--whitelist: Nama topik sumber. Untuk menentukan beberapa topik, pisahkan dengan tanda pipa (|), misalnyatopicA|topicB|topicC. -
--topic.mapping.file: File konfigurasi untuk pemetaan topik. -
KAFKA_LOG4J_OPTS: Path ke file konfigurasi 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 & -
-
Periksa log untuk mencari error. Output berikut menunjukkan startup yang berhasil:
-
[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
-
-
Login ke Konsol DataHub untuk memverifikasi bahwa data telah ditulis ke DataHub.
Pada halaman Data Bus, buka Topik tujuan dalam Proyek tujuan, misalnya
kafkatest/test. Pada tab Shard List, pastikan status shard adalahACTIVEdan terdapat waktu data terbaru yang ditampilkan. Klik Sample di kolom Actions. Pada panel yang muncul, pilih Shard ID, tentukan jumlah catatan yang akan diambil, lalu klik Sample. Pada tabel pratinjau data, verifikasi bahwa data berhasil ditulis ke DataHub.