All Products
Search
Document Center

DataHub:Mode kompatibilitas Kafka DataHub

Last Updated:Jun 26, 2026

Mode kompatibilitas Kafka DataHub

DataHub kini kompatibel dengan protokol Apache Kafka. Anda dapat menggunakan SDK Kafka standar untuk terhubung ke DataHub guna memublikasikan dan berlangganan data.

Pemetaan konsep antara DataHub dan Kafka

Kafka

DataHub

topic

Project.topic

partition

shard

offset

sequence

Kafka topic

DataHub juga memiliki topic. Namun, DataHub mencakup lapisan resource tambahan yang disebut project, yang tidak ada di Kafka. Untuk memastikan kompatibilitas, sebuah topic Kafka berkorespondensi dengan kombinasi project dan topic DataHub, yang digabungkan dengan tanda titik (.). Sebagai contoh, jika Anda memiliki project bernama test_project dan topic bernama test_topic di DataHub, maka nama topic Kafka yang sesuai adalah test_project.test_topic.

Kafka partition

Sebuah partition Kafka secara langsung setara dengan shard DataHub. Keduanya merepresentasikan antrian data terurut.

Kafka consumer group

Kelompok konsumen DataHub berperilaku mirip dengan kelompok konsumen Kafka. Anda dapat membuat kelompok konsumen dalam suatu project dan mengaitkannya dengan topic yang ingin Anda langgani. Hal ini memungkinkan kelompok tersebut untuk berlangganan ke beberapa topic dalam project tersebut. Jika kelompok konsumen dikaitkan ke suatu topic, langganan akan dibuat secara otomatis, yang dapat Anda lihat pada halaman daftar langganan topic tersebut. Menghapus langganan ini mencegah kelompok tersebut berlangganan ke topic dan menghapus semua offset konsumen sebelumnya.

Karena kelompok konsumen merupakan sub-resource dari project, Anda harus menentukannya bersama dengan project. Sebagai contoh, jika project DataHub Anda adalah test_project dan kelompok konsumen Anda adalah test_group, maka nama grup Kafka-nya adalah test_project.test_group.

Pada kotak dialog Create consumer group, masukkan Name dan Description, pilih topic yang ingin dilangganan di area Topic, klik tombol > untuk menambahkannya ke daftar terpilih di sebelah kanan, lalu klik Create.

Setiap kelompok dapat berlangganan hingga maksimal 50 topic. Jika Anda perlu berlangganan lebih banyak, harap ajukan tiket.

Kafka Record

Sebuah record Kafka menggunakan format pasangan kunci-nilai. Sebaliknya, record DataHub tersedia dalam dua format: Tuple dan Blob. Record Tuple berisi data terstruktur kuat, sedangkan record Blob berisi data biner. Secara umum, sebuah record Kafka terdiri dari dua bagian: Header dan pasangan kunci-nilai.

Kafka Header

Header Kafka dapat menambahkan informasi tambahan ke data, yang memiliki fungsi yang sama dengan atribut DataHub. Jika Anda menggunakan client Kafka untuk menulis data yang menyertakan header, informasi header tersebut akan disimpan sebagai atribut DataHub. Jika nilai header Kafka adalah NULL, header tersebut akan diabaikan. Jangan gunakan "__kafka_key__" sebagai kunci header karena merupakan kunci internal yang dicadangkan.

Data pasangan kunci-nilai Kafka

  • Jika topic DataHub bertipe Tuple, kuncinya berkorespondensi dengan kolom String pertama dan nilainya berkorespondensi dengan kolom String kedua.

  • Jika topic DataHub bertipe Blob, nilainya berkorespondensi dengan konten data, dan kuncinya ditempatkan dalam atribut dengan format <"__kafka_key__", key>.

Kafka offset

Offset Kafka adalah bilangan bulat 64-bit. Untuk suatu partition, offset dimulai dari 0 dan bertambah secara otomatis. Hal ini memastikan bahwa setiap catatan data memiliki offset unik, yang terutama digunakan untuk mencatat offset konsumen. Sequence di DataHub memiliki makna yang sama dengan offset di Kafka. Oleh karena itu, Anda dapat menganggap sequence DataHub secara langsung setara dengan offset Kafka.

Batasan

DataHub tidak mendukung transaksi Kafka, idempotensi, SchemaRegistry, atau Log Compaction.

Memulai Cepat

  1. Gunakan versi client Kafka yang direkomendasikan, yaitu 2.4.0. Versi ini kompatibel dengan versi 0.10.0 hingga 4.0.

  2. Buat resource Project, Topic, dan Group yang sesuai di DataHub.

  3. Ubah metode autentikasi keamanan menjadi SASL/SSL, gunakan mekanisme PLAIN untuk SASL, dan konfigurasikan Pasangan Kunci Akses Alibaba Cloud (AK/SK) Anda.

  4. Ubah informasi broker Kafka menjadi titik akhir DataHub (lihat daftar domain layanan).

Pembuatan resource

Pertama, login ke Konsol DataHub dan buat project.

Pada kotak dialog New Project, atur Name menjadi test_kafka_project dan Description menjadi test kafka.

Buat topic

Pada halaman New Topic di Konsol DataHub, pilih Direct Create sebagai metode pembuatan, masukkan test_kafka sebagai nama, dan pilih TUPLE sebagai tipe. Di bagian detail Skema, tambahkan dua bidang bernama key dan value, atur tipe keduanya menjadi STRING, dan pilih Allow NULL. Atur Shard Count menjadi 1 dan Lifecycle menjadi 3, aktifkan Shard Scaling Mode, dan nonaktifkan Multi-version. Masukkan deskripsi dan klik Create.

Buat kelompok dan kaitkan topic yang perlu dikonsumsi. Anda juga dapat mengubah daftar topic yang dikaitkan setelah kelompok dibuat. Jika Anda tidak perlu mengonsumsi pesan, Anda dapat melewati langkah ini.

Pada panel New Group, masukkan Name dan Description, gunakan daftar transfer untuk memindahkan topic yang akan dikonsumsi dari daftar tersedia di kiri ke daftar terpilih di kanan, lalu klik Create.

Metode Autentikasi

Buat file bernama kafka_client_producer_jaas.conf dan simpan ke path apa pun. Isi file tersebut adalah sebagai berikut.

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="accessId"
  password="accessKey";
};

Buat file bernama kafka_client_producer_jaas.conf dan simpan ke direktori apa pun. File tersebut harus berisi konten berikut:

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="YOUR_ACCESS_ID"
  password="YOUR_ACCESS_KEY";
};

Contoh Producer

File pom client kafka

<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-clients</artifactId>
  <version>2.4.0</version>
</dependency>

File pom.xml client Kafka:

<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-clients</artifactId>
  <version>2.4.0</version>
</dependency>

ProducerExample

Contoh Konsumen

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());
            }
        }
    }
}

Contoh Streams

Kode ini membaca data dari input di test_project, mengonversi string kunci dan nilai menjadi huruf kecil, lalu menulis ulang data tersebut ke 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));
        }
    }
}

Contoh ini membaca data dari topic input dalam project test_project, mengonversi string kunci dan nilai menjadi huruf kecil, lalu menulis hasilnya ke topic output.

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));
        }
    }
}

Setelah Anda menjalankan job Streams, proses penugasan shard membutuhkan waktu sekitar satu menit. Anda kemudian dapat melihat jumlah tugas saat ini di konsol. Jumlah tugas tersebut sesuai dengan jumlah shard dalam topic input. Dalam contoh ini, topic input memiliki tiga shard.

Hasil pengambilan sampel output menunjukkan bahwa data telah ditulis dengan benar. Setelah diproses oleh TestMapper dalam tugas Streams, input huruf kapital dikonversi menjadi output huruf kecil: kuncinya adalah aaaa, cccc, dan eeee, serta nilai yang sesuai adalah bbbb, dddd, dan ffff.

Migrasi data dari Kafka self-managed ke DataHub

  1. Alihkan alamat broker. Untuk informasi lebih lanjut, lihat daftar nama domain di bawah.

  2. Buat resource di DataHub dan modifikasi nama resource dalam kode Anda. Untuk informasi lebih lanjut, lihat bagian pembuatan resource di atas.

  3. Alihkan metode autentikasi ke SASL_SSL. Mekanisme autentikasinya adalah PLAIN. Atur username ke ID AccessKey Alibaba Cloud (AK) Anda dan password ke Secret AccessKey (SK) Anda.

Lampiran

Ikhtisar Konfigurasi

C=Consumer, P=Producer, S=Streams

Parameter

C/P/S

Nilai yang direkomendasikan

Wajib

Deskripsi

bootstrap.servers

*

Lihat daftar nama domain.

Ya

security.protocol

*

SASL_SSL

Ya

Untuk memastikan transmisi data aman, penulisan data dari Kafka ke DataHub menggunakan enkripsi SSL secara default.

sasl.mechanism

*

PLAIN

Ya

Otentikasi AccessKey. Hanya PLAIN yang didukung.

compression.type

P

LZ4

Tidak

Menentukan apakah kompresi diaktifkan untuk transmisi data. Saat ini, hanya LZ4 yang didukung.

enable.idempotence

P

false

Tidak

Client Kafka versi 3.0.1 dan yang lebih baru mengaktifkan idempotensi secara default. Karena DataHub tidak mendukung idempotensi, Anda harus menonaktifkan fitur ini secara manual. Konfigurasi ini tidak diperlukan untuk client versi sebelum 3.0.1.

group.id

C

project.group

Ya

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

Tidak

Nilai default untuk Kafka adalah RangeAssignor. DataHub saat ini hanya mendukung RangeAssignor. Jangan ubah konfigurasi ini.

session.timeout.ms

C/S

[60000, 180000]

Tidak

Nilai default di Kafka adalah 10000. Namun, DataHub menerapkan batas minimum 60000, sehingga nilai default-nya menjadi 60000.

heartbeat.interval.ms

C/S

Dua pertiga dari nilai session.timeout.ms.

Tidak

Nilai default di Kafka adalah 3000. Karena session.timeout.ms default-nya 60000, atur eksplisit nilai ini menjadi 40000 untuk mencegah permintaan heartbeat yang terlalu sering.

application.id

S

project.topic:subId atau project.group

Ya

Jika Anda menggunakan project.topic:subId, nilai tersebut harus sesuai dengan topic yang dilangganan. Jika tidak, data tidak dapat dibaca. Kami merekomendasikan penggunaan project.group.

Daftar Nama Domain Layanan

Nama wilayah

Wilayah

Titik akhir publik

Titik akhir VPC ECS

China Timur 1 (Hangzhou)

cn-hangzhou

dh-cn-hangzhou.aliyuncs.com:9092

dh-cn-hangzhou-int-vpc.aliyuncs.com:9094

China Timur 2 (Shanghai)

cn-shanghai

dh-cn-shanghai.aliyuncs.com:9092

dh-cn-shanghai-int-vpc.aliyuncs.com:9094

China Utara 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 Selatan 1 (Shenzhen)

cn-shenzhen

dh-cn-shenzhen.aliyuncs.com:9092

dh-cn-shenzhen-int-vpc.aliyuncs.com:9094

China Utara 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

Eropa Tengah 1 (Frankfurt)

eu-central-1

dh-eu-central-1.aliyuncs.com:9092

dh-eu-central-1-int-vpc.aliyuncs.com:9094

China Timur 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

Kompatibilitas API Kafka

Dokumentasi resmi Kafka mencantumkan semua API yang tersedia. Untuk membantu Anda memahami kompatibilitas KOD, kami juga menyediakan daftar API yang kompatibel.

API

Deskripsi

Produce

Menulis data.

Fetch

Membaca data.

ListOffsets

Mendapatkan offset berdasarkan waktu. Versi sebelumnya mengembalikan daftar offset, sedangkan versi terbaru mengembalikan satu offset.

Metadata

Mendapatkan metadata untuk operasi baca dan tulis.

OffsetCommit

Terapkan offset.

OffsetFetch

Mendapatkan offset.

FindCoordinator

Menemukan broker yang menjadi host coordinator dan mengembalikan alamat IP virtual (VIP)-nya.

JoinGroup

Bergabung ke kelompok.

Heartbeat

Mengirim heartbeat.

LeaveGroup

Keluar dari kelompok.

SyncGroup

Digunakan oleh pemimpin kelompok untuk mengirim rencana penugasan partisi, dan oleh semua anggota untuk menerimanya.

SaslHandshake

Menangani handshake Simple Authentication and Security Layer (SASL) untuk autentikasi.

ApiVersions

Mendapatkan semua API yang tersedia.

CreateTopics

Membuat topic.

DeleteTopics

Menghapus topic.

OffsetForLeaderEpoch

Mendapatkan offset terbaru.

SaslAuthenticate

Mengotentikasi client menggunakan SASL.

CreatePartitions

Menambahkan partisi.

DeleteGroups

Menghapus kelompok.

OffsetDelete

Menghapus offset.