DataHub sepenuhnya kompatibel dengan protokol Apache Kafka. Anda dapat menggunakan klien Kafka native untuk membaca data dari dan menulis data ke DataHub.
Pemetaan Kafka ke DataHub
Jenis Topik
Kafka dan DataHub memiliki mekanisme penskalaan topik yang berbeda. Untuk memastikan kompatibilitas dengan perilaku Kafka, Anda harus mengatur mode penskalaan ke ONLY_EXTEND saat membuat topik DataHub. Dalam mode ini, Anda hanya dapat menambahkan shard baru ke suatu topik. Mode ini tidak mengizinkan pemisahan atau penggabungan shard, serta belum mendukung penghapusan shard.
Penamaan Topik
Nama topik Kafka dipetakan ke proyek dan topik DataHub, dipisahkan oleh titik (.). Pemetaan ini mengikuti aturan berikut:
Bagian sebelum titik
.pertama adalah proyek DataHub, dan bagian setelahnya adalah topik DataHub. Misalnya,test_project.test_topicdipetakan ke proyektest_projectdan topiktest_topic.Jika nama berisi beberapa karakter
., hanya titik.pertama yang berfungsi sebagai pemisah. Semua karakter.dan-berikutnya diganti dengan_.
Partisi
Setiap shard aktif di DataHub bersesuaian dengan satu partisi di Kafka. Misalnya, jika sebuah topik memiliki lima shard aktif, maka itu setara dengan topik Kafka dengan lima partisi. Saat menulis data, Anda dapat menentukan ID partisi dalam rentang [0, 4]. Jika tidak menentukan partisi, klien Kafka akan secara otomatis menetapkan salah satunya.
Topik Tuple
Saat menulis data dari Kafka ke topik Tuple, skema topik harus memiliki satu atau dua kolom, keduanya bertipe STRING. Jika tidak, operasi penulisan akan gagal.
Jika skema memiliki satu kolom, hanya nilai (value) yang ditulis, sedangkan kunci (key) dibuang.
Jika skema memiliki dua kolom, kolom pertama dan kedua masing-masing bersesuaian dengan kunci dan nilai.
Selain itu, jangan menulis data biner ke topik Tuple karena akan menyebabkan karakter acak (garbled). Untuk menyimpan data biner, gunakan topik Blob.
Topik Blob
Saat menulis data dari Kafka ke topik Blob, nilai pesan Kafka ditulis ke bidang Blob. Jika kunci pesan tidak NULL, kunci tersebut ditulis sebagai atribut DataHub dengan nama atribut __kafka_key__ dan nilainya adalah kunci pesan Kafka.
Header
Header Kafka dipetakan ke atribut DataHub. Jika nilai header bernilai NULL, header tersebut diabaikan dan tidak ditulis sebagai atribut. Kami menyarankan agar Anda tidak menggunakan __kafka_key__ sebagai kunci header untuk menghindari konflik dengan nama atribut bawaan pada topik Blob.
Kelompok Konsumen
Di DataHub, ID langganan berperan sebagai kelompok konsumen tetapi hanya dapat berlangganan ke satu topik saja. Sebaliknya, kelompok konsumen Kafka dapat berlangganan ke beberapa topik secara bersamaan. Untuk memberikan kompatibilitas dengan model langganan Kafka, DataHub menyediakan fitur kelompok. Anda dapat membuat kelompok dalam suatu proyek dan mengaitkannya ke beberapa topik, sehingga memungkinkan Anda berlangganan ke semuanya dalam satu kelompok.
Kelompok tersebut secara internal mengelola beberapa langganan DataHub di sisi server. Setelah Anda mengaitkan topik, kelompok tersebut secara otomatis membuat langganan, yang muncul di daftar langganan pada halaman detail topik. Jangan menghapus langganan ini secara manual. Melakukannya akan mencegah kelompok berlangganan ke topik tersebut dan menyebabkan semua offset konsumsi yang ada hilang.
Satu kelompok dapat berlangganan hingga maksimal 50 topik. Untuk berlangganan lebih banyak topik, ajukan tiket.
Parameter Kafka
C=Consumer, P=Producer, S=Streams
Parameter | C/P/S | Nilai | Wajib | Deskripsi |
bootstrap.servers | * | Lihat bagian Kafka endpoints. | Ya | |
security.protocol | * | SASL_SSL | Ya | Untuk memastikan keamanan data, koneksi dari Kafka ke DataHub menggunakan enkripsi SSL secara default. |
sasl.mechanism | * | PLAIN | Ya | Mekanisme autentikasi untuk kredensial AccessKey. Hanya PLAIN yang didukung. |
compression.type | P | LZ4 | Tidak | Menentukan jenis kompresi untuk pesan. Saat ini, hanya LZ4 yang didukung. |
group.id | C | project.topic:subId atau project.group | Ya | Jika Anda menggunakan format |
partition.assignment.strategy | C | org.apache.kafka.clients.consumer.RangeAssignor | Tidak | Strategi penetapan partisi default di Kafka adalah |
session.timeout.ms | C/S | [60000, 180000] | Tidak | Nilai default di Kafka adalah 10.000 ms. Namun, karena DataHub memerlukan minimal 60.000 ms, nilai ini secara otomatis disesuaikan menjadi 60.000 ms. |
heartbeat.interval.ms | C/S | Disarankan: 2/3 dari session.timeout.ms | Tidak | Nilai default Kafka adalah 3.000 ms. Karena |
application.id | S | project.topic:subId atau project.group | Ya | Jika Anda menggunakan format |
Tabel ini mencantumkan parameter utama yang perlu ditinjau saat menggunakan klien Kafka dengan DataHub. Parameter sisi klien lainnya, seperti retries dan batch.size, berperilaku sebagaimana dalam Kafka native. Parameter sisi server tidak mengubah perilaku aktual DataHub. Misalnya, terlepas dari nilai acks, DataHub hanya mengembalikan konfirmasi setelah data benar-benar ditulis.
Endpoint Kafka
Wilayah | ID Wilayah | Endpoint publik | Titik akhir ECS (jaringan klasik) | Endpoint ECS (VPC) |
Tiongkok (Hangzhou) | cn-hangzhou | dh-cn-hangzhou.aliyuncs.com:9092 | dh-cn-hangzhou.aliyun-inc.com:9093 | dh-cn-hangzhou-int-vpc.aliyuncs.com:9094 |
Tiongkok (Shanghai) | cn-shanghai | dh-cn-shanghai.aliyuncs.com:9092 | dh-cn-shanghai.aliyun-inc.com:9093 | dh-cn-shanghai-int-vpc.aliyuncs.com:9094 |
Tiongkok (Beijing) | cn-beijing | dh-cn-beijing.aliyuncs.com:9092 | dh-cn-beijing.aliyun-inc.com:9093 | dh-cn-beijing-int-vpc.aliyuncs.com:9094 |
Tiongkok (Zhangjiakou) | cn-zhangjiakou | dh-cn-zhangjiakou.aliyuncs.com:9092 | dh-cn-zhangjiakou.aliyun-inc.com:9093 | dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094 |
Tiongkok (Shenzhen) | cn-shenzhen | dh-cn-shenzhen.aliyuncs.com:9092 | dh-cn-shenzhen.aliyun-inc.com:9093 | dh-cn-shenzhen-int-vpc.aliyuncs.com:9094 |
Singapura | ap-southeast-1 | dh-ap-southeast-1.aliyuncs.com:9092 | dh-ap-southeast-1.aliyun-inc.com:9093 | dh-ap-southeast-1-int-vpc.aliyuncs.com:9094 |
Malaysia (Kuala Lumpur) | ap-southeast-3 | dh-ap-southeast-3.aliyuncs.com:9092 | dh-ap-southeast-3.aliyun-inc.com:9093 | dh-ap-southeast-3-int-vpc.aliyuncs.com:9094 |
Jerman (Frankfurt) | eu-central-1 | dh-eu-central-1.aliyuncs.com:9092 | dh-eu-central-1.aliyun-inc.com:9093 | dh-eu-central-1-int-vpc.aliyuncs.com:9094 |
Tiongkok Timur 2 Finance | cn-shanghai-finance-1 | dh-cn-shanghai-finance-1.aliyuncs.com:9092 | dh-cn-shanghai-finance-1.aliyun-inc.com:9093 | dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094 |
Tiongkok (Hong Kong) | cn-hongkong | dh-cn-hongkong.aliyuncs.com:9092 | dh-cn-hongkong.aliyun-inc.com:9093 | dh-cn-hongkong-int-vpc.aliyuncs.com:9094 |
Buat topik
Buat topik di Konsol
Saat membuat topik, aktifkan Shard Expand Mode.
Buat topik menggunakan SDK
Anda tidak dapat membuat topik menggunakan API Kafka. Anda harus menggunakan SDK DataHub dan mengatur
ExpandModekeONLY_EXTEND. Versi dependensi Maven yang diperlukan adalah 2.19.0 atau lebih baru.Kami menyarankan agar Anda menggunakan variabel lingkungan untuk mengonfigurasi ID AccessKey dan AccessKey Secret Anda serta menghindari hard-coding ke dalam kode proyek. Pasangan Kunci Akses Akun Alibaba Cloud memiliki izin untuk semua operasi API. Untuk keamanan yang lebih baik, gunakan pasangan Kunci Akses dari Pengguna RAM untuk akses API atau operasi harian untuk mengurangi risiko kebocoran kredensial.
datahub.endpoint=<yourEndpoint> datahub.accessId=<yourAccessKeyId> datahub.accessKey=<yourAccessKeySecret><dependency> <groupId>com.aliyun.datahub</groupId> <artifactId>aliyun-sdk-datahub</artifactId> <version>2.19.0-public</version> </dependency>@Value("${datahub.endpoint}") String endpoint ; @Value("${datahub.accessId}") String accessId; @Value("${datahub.accessKey}") String accessKey; public class CreateTopic { public static void main(String[] args) { DatahubClient datahubClient = DatahubClientBuilder.newBuilder() .setDatahubConfig( new DatahubConfig(endpoint, new AliyunAccount(accessId, accessKey))) .build(); int shardCount = 1; int lifeCycle = 7; try { datahubClient.createTopic("test_project", "test_topic", shardCount, lifeCycle, RecordType.BLOB, "comment", ExpandMode.ONLY_EXTEND); } catch (DatahubClientException e) { e.printStackTrace(); } } }
Buat kelompok
Buat kelompok di konsol
Klik Create Group, lalu tambahkan topik yang ingin Anda langgani dari daftar di sebelah kanan. Anda dapat mengubah topik yang dikaitkan setelah kelompok dibuat. Kelompok tersebut secara otomatis membuat langganan, yang muncul di halaman daftar langganan topik.
Buat kelompok menggunakan SDK
Versi dependensi Maven harus 2.21.6-public atau lebih baru.
<dependency> <groupId>com.aliyun.datahub</groupId> <artifactId>aliyun-sdk-datahub</artifactId> <version>2.21.6-public</version> </dependency>@Value("${datahub.endpoint}") String endpoint ; @Value("${datahub.accessId}") String accessId; @Value("${datahub.accessKey}") String accessKey; public class CreateGroup { public static void main(String[] args) { DatahubClient datahubClient = DatahubClientBuilder.newBuilder() .setDatahubConfig( new DatahubConfig(endpoint, new AliyunAccount(accessId, accessKey))) .build(); List<String> topicList = new ArrayList<>(); topicList.add("test_project.topic1"); topicList.add("test_project.topic2"); topicList.add("test_project.topic3"); try { // Buat kelompok Kafka. datahubClient.createKafkaGroup("test_project", "test_topic", "test comment"); // Kaitkan topik ke kelompok untuk langganan. datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD); } catch (DatahubClientException e) { e.printStackTrace(); } } }
Contoh Produsen
File kafka_client_producer_jaas.conf
Buat file bernama kafka_client_producer_jaas.conf di direktori mana pun dan tambahkan konten berikut.
KafkaClient {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="yourAccessKeyId"
password="yourAccessKeySecret";
};Dependensi Maven
Versi klien Kafka harus 0.10.0.0 atau lebih baru. Kami menyarankan versi 2.4.0.
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>Kode contoh
public class ProducerExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/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("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
properties.put("compression.type", "lz4");
String KafkaTopicName = "test_project.test_topic";
Producer<String, String> producer = new KafkaProducer<String, String>(properties);
try {
List<Header> headers = new ArrayList<>();
RecordHeader header1 = new RecordHeader("key1", "value1".getBytes());
RecordHeader header2 = new RecordHeader("key2", "value2".getBytes());
headers.add(header1);
headers.add(header2);
ProducerRecord<String, String> record = new ProducerRecord<>(KafkaTopicName, 0, "key", "Hello DataHub!", headers);
// Kirim sinkron
producer.send(record).get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
} finally {
producer.close();
}
}
}Hasil
Setelah kode berhasil dijalankan, Anda dapat mengambil data sampel untuk memverifikasi hasilnya.
Contoh Consumer
Untuk informasi tentang cara membuat filekafka_client_producer_jaas.conf dan menambahkan dependensi Maven, lihat contoh producer.
Saat consumer baru bergabung, penugasan shard memerlukan waktu 10 hingga 20 detik. Setelah penugasan selesai, consumer dapat mulai mengonsumsi data.
Kode contoh
Menggunakan kelompok Kafka (disarankan)
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", "src/main/resources/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");
// Atur group.id ke format project.group.
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");
// Dengan menggunakan kelompok Kafka, Anda dapat berlangganan ke beberapa topik.
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());
}
}
}
}Menggunakan project.topic:subId
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.Collections;
import java.util.Properties;
public class ConsumerExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/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");
// Atur group.id ke format project.topic:subId.
properties.put("group.id", "test_project.test_topic:1611039998153N71KM");
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);
// Saat menggunakan format project.topic:subId, Anda hanya dapat berlangganan ke satu topik.
kafkaConsumer.subscribe(Collections.singletonList("test_project.test_topic"));
while (true) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.toString());
}
}
}
}Hasil
Setelah kode berhasil dijalankan, Anda dapat melihat data yang dikonsumsi di terminal Anda.
ConsumerRecord(topic = test_project.test_topic, partition = 0, leaderEpoch = 0, offset = 0, LogAppendTime = 1611040892661, serialized key size = 3, serialized value size = 14, headers = RecordHeaders(headers = [RecordHeader(key = key1, value = [118, 97, 108, 117, 101, 49]), RecordHeader(key = key2, value = [118, 97, 108, 117, 101, 50])], isReadOnly = false), key = key, value = Hello DataHub!)Dalam contoh ini, semua catatan data yang dikembalikan dalam satu permintaan memiliki LogAppendTime yang sama, yaitu timestamp terbaru di antara semua catatan dalam batch tersebut.
Contoh Streams
Dependensi Maven
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>2.4.0</version>
</dependency>Kode contoh
Contoh ini membaca data dari topik input dalam test_project, mengonversi string kunci dan nilai menjadi huruf kecil, lalu menulis hasilnya ke topik output.
public class StreamExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/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.input:1611293595417QH0WL");
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));
}
}
}Hasil
Setelah Anda memulai tugas Streams, penugasan shard memerlukan waktu sekitar satu menit. Setelah itu, Anda dapat melihat jumlah tugas saat ini di Konsol. Jumlah tugas tersebut sesuai dengan jumlah shard di topik input. Dalam contoh ini, topik input memiliki tiga shard.
currently assigned active tasks: [0_0, 0_1, 0_2]
currently assigned standby tasks: []
revoked active tasks: []
revoked standby tasks: []Setelah shard ditugaskan, Anda dapat menulis data uji seperti (AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF) ke topik input. Kemudian, ambil data sampel dari topik output untuk memverifikasi bahwa data telah ditulis dengan benar.
Catatan penggunaan
Transaksi dan idempotensi tidak didukung.
Klien Kafka tidak dapat membuat topik secara otomatis di DataHub. Anda harus membuat topik sebelum menulis data ke dalamnya.
Saat menggunakan ID langganan (
project.topic:subid) sebagaigroup.id, consumer hanya dapat berlangganan ke satu topik. Untuk berlangganan ke beberapa topik, gunakan kelompok DataHub.Timestamp untuk data yang dibaca oleh consumer selalu merupakan LogAppendTime, yang menunjukkan kapan data ditulis ke DataHub. Semua catatan dalam satu permintaan fetch memiliki timestamp yang sama: timestamp terbaru dalam batch tersebut. Artinya, timestamp pembacaan mungkin lebih lambat daripada waktu penulisan sebenarnya.
Aplikasi Streams hanya mendukung satu topik input tetapi dapat memiliki beberapa topik output.
Hanya tugas Streams tanpa status (stateless) yang didukung.
Versi Kafka yang didukung berkisar dari 0.10.0 hingga 2.4.0.
FAQ
Koneksi terputus saat menulis data
Selector - [Producer clientId=producer-1] Connection with dh-cn-shenzhen.aliyuncs.com disconnected
java.io.EOFException
at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:573)
...Permintaan metadata Kafka dan permintaan penulisan data menggunakan koneksi yang berbeda.
Klien pertama kali membuat koneksi untuk mengambil metadata. Lalu, klien menggunakan informasi broker yang dikembalikan untuk membuat koneksi kedua untuk menulis data. Semua permintaan selanjutnya dikirim melalui koneksi kedua ini.
Koneksi pertama, yang kini menganggur, secara otomatis ditutup oleh server setelah timeout. Hal ini dapat menghasilkan error pemutusan koneksi di log. Anda dapat mengabaikan error ini jika data berhasil ditulis.
Klien Kafka gagal dimulai
Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed
Caused by: javax.net.ssl.SSLHandshakeException: No subject alternative names matching IP address 100.67.134.161 foundTambahkan properti berikut ke konfigurasi Anda: properties.put("ssl.endpoint.identification.algorithm", "");.
DisconnectException saat konsumsi
[INFO][Consumer clientId=client-id, groupId=consumer-project.topic:subid] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: {}.
org.apache.kafka.common.errors.DisconnectExceptionKlien Kafka harus mempertahankan koneksi TCP persisten dengan server. Pengecualian ini biasanya disebabkan oleh fluktuasi jaringan. Klien memiliki logika retry bawaan, sehingga error ini umumnya tidak memengaruhi konsumsi.