All Products
Search
Document Center

DataHub:Migrasikan data ke DataHub dengan Kafka Mirror Maker

Last Updated:Jun 26, 2026

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.

Catatan
  • 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

  1. Unggah paket kafka_mirror_datahub.tgz ke server Kafka sumber Anda dan ekstrak paket tersebut.

    tar -zxvf kafka_mirror_datahub.tgz
  1. Buka direktori config untuk mengubah file konfigurasi sumber dan tujuan. File-file tersebut dijelaskan sebagai berikut:

    1. 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.StringDeserializer

    b. 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=false

    Penjelasan parameter

    • Untuk daftar nama domain bootstrap.servers tujuan, lihat Kompatibilitas Kafka. Titik akhir pada contoh ini, dh-cn-zhangjiakou-pre.aliyuncs.com:9092, sesuai dengan wilayah China (Zhangjiakou).

    • Parameter sasl.jaas.config menentukan modul login dan kredensial yang diperlukan untuk otentikasi SASL. Ganti AccessKey ID dan AccessKey Secret dengan Informasi AccessKey Anda.

    • Untuk informasi lebih lanjut tentang item konfigurasi, lihat Kompatibilitas Kafka.

  1. Konfigurasikan file topic-map.properties untuk 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
  1. Konfigurasikan file log4j.properties untuk logging.

    1. Buat file log4j.properties.

    2. Gunakan templat berikut untuk konfigurasi Anda:

      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. 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 (|), misalnya topicA|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 &
  3. Periksa log untuk mencari error. Output berikut menunjukkan startup yang berhasil:

    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. 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 adalah ACTIVE dan 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.