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
-
Gunakan versi client Kafka yang direkomendasikan, yaitu 2.4.0. Versi ini kompatibel dengan versi 0.10.0 hingga 4.0.
-
Buat resource Project, Topic, dan Group yang sesuai di DataHub.
-
Ubah metode autentikasi keamanan menjadi SASL/SSL, gunakan mekanisme PLAIN untuk SASL, dan konfigurasikan Pasangan Kunci Akses Alibaba Cloud (AK/SK) Anda.
-
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
-
Alihkan alamat broker. Untuk informasi lebih lanjut, lihat daftar nama domain di bawah.
-
Buat resource di DataHub dan modifikasi nama resource dalam kode Anda. Untuk informasi lebih lanjut, lihat bagian pembuatan resource di atas.
-
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 |
|
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. |