Hubungkan aplikasi Anda ke instans ApsaraMQ for Kafka melalui titik akhir default, SSL, atau SASL untuk mengirim dan menerima pesan.
Prasyarat
-
Instal JDK 1.8 atau versi yang lebih baru. Untuk informasi selengkapnya, lihat Install JDK.
-
Instal Maven 2.5 atau versi yang lebih baru. Untuk informasi selengkapnya, lihat Install Maven.
-
Instal tool build.
Topik ini menggunakan IntelliJ IDEA Ultimate sebagai contoh.
-
Beli dan deploy instans ApsaraMQ for Kafka.
-
Instans VPC: Hanya menyediakan titik akhir default yang dapat diakses dari dalam VPC yang sama.
-
Instans Internet/VPC: Menyediakan titik akhir default dan titik akhir SSL, serta dapat diakses melalui Internet atau dari VPC.
Catatan-
Anda dapat mengonversi instans VPC menjadi instans Internet/VPC dan sebaliknya. Untuk informasi selengkapnya, lihat Upgrade atau downgrade konfigurasi instans.
-
Titik akhir SASL dinonaktifkan secara default dan tidak muncul di Konsol. Untuk menggunakan titik akhir SASL, Anda harus mengaktifkannya secara manual. Untuk informasi selengkapnya, lihat Berikan izin kepada pengguna SASL.
-
Untuk informasi selengkapnya tentang kasus penggunaan masing-masing titik akhir, lihat Perbandingan titik akhir.
-
Instal dependensi Java
Dependensi berikut diperlukan. Dependensi ini sudah termasuk dalam file pom.xml di folder kafka-java-demo, sehingga Anda tidak perlu menambahkannya secara manual.
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.6.0</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>1.7.6</version>
</dependency>
Kami menyarankan agar Anda menjaga konsistensi versi utama library client dengan versi utama instans ApsaraMQ for Kafka Anda. Anda dapat menemukan versi utama instans ApsaraMQ for Kafka Anda di halaman Instance Details di Konsol ApsaraMQ for Kafka.
Siapkan konfigurasi
-
Opsional: Unduh sertifikat root SSL. Jika Anda menggunakan titik akhir SSL, Anda harus mengunduh sertifikat ini.
-
Buka Aliware-kafka-demos, klik
, lalu unduh dan ekstrak proyek demo. -
Impor folder kafka-java-demo dari proyek yang telah diekstrak ke IntelliJ IDEA.
-
Opsional:Jika Anda menggunakan titik akhir SSL atau titik akhir SASL untuk menghubungkan ke instans, Anda harus memodifikasi file kafka_client_jaas.conf. Untuk informasi selengkapnya tentang titik akhir instans, lihat Perbandingan titik akhir.
KafkaClient { org.apache.kafka.common.security.plain.PlainLoginModule required username="xxxx" password="xxxx"; };Untuk instans VPC, hanya sumber daya dalam VPC yang sama yang dapat mengakses instans Kafka, sehingga memastikan transmisi data yang aman dan privat. Untuk keamanan yang lebih tinggi, Anda dapat mengaktifkan fitur ACL sehingga pesan hanya ditransmisikan setelah autentikasi SASL. Anda dapat menggunakan mekanisme PLAIN atau SCRAM sesuai kebutuhan keamanan Anda. Untuk informasi selengkapnya, lihat Aktifkan fitur ACL.
Untuk instans Internet/VPC, pesan yang ditransmisikan melalui Internet harus diautentikasi dan dienkripsi. SASL dengan mekanisme PLAIN harus digunakan di atas lapisan transport SSL. Protokol SASL_SSL mencegah transmisi teks biasa melalui Internet.
username dan password pada contoh tersebut adalah username dan password SASL untuk instans tersebut.
-
Jika fitur ACL dinonaktifkan untuk instans yang dapat diakses melalui Internet, Anda dapat memperoleh username dan password pengguna default dari bagian Configuration Information di halaman Instance Details di ApsaraMQ for Kafka console.
-
Jika fitur ACL diaktifkan untuk instans tersebut, pastikan pengguna SASL bertipe PLAIN dan telah diberikan izin untuk mengirim serta menerima pesan. Untuk informasi selengkapnya, lihat Berikan izin kepada pengguna SASL.
-
-
Modifikasi file konfigurasi kafka.properties.
##==============================Common parameters============================== bootstrap.servers=xxxxxxxxxxxxxxxxxxxxx topic=xxx group.id=xxx ##=======================Configure the following parameters based on your requirements======================== ##SSL endpoint configuration ssl.truststore.location=/xxxx/only.4096.client.truststore.jks ##The value of ssl.truststore.password must be KafkaOnsClient and cannot be changed. ssl.truststore.password=KafkaOnsClient ##Hostname verification algorithm. Keep this parameter empty. Do not change the setting. ssl.endpoint.identification.algorithm= java.security.auth.login.config=/xxxx/kafka_client_jaas.conf ##SASL endpoint with the PLAIN mechanism java.security.auth.login.config.plain=/xxxx/kafka_client_jaas_plain.conf ##SASL endpoint with the SCRAM mechanism java.security.auth.login.config.scram=/xxxx/kafka_client_jaas_scram.confParameter
Deskripsi
bootstrap.servers
Titik akhir. Anda dapat memperoleh titik akhir dari bagian Endpoint Information di halaman Instance Details di ApsaraMQ for Kafka console.
topic
Nama topik. Anda dapat memperoleh nama topik di halaman Topics di ApsaraMQ for Kafka console.
group.id
Group dari instans. Anda dapat memperolehnya di halaman Groups di ApsaraMQ for Kafka console.
CatatanParameter ini opsional untuk produsen tetapi wajib untuk konsumen.
ssl.truststore.location
Jalur lokal sertifikat root SSL yang telah diunduh. Ganti xxxx dengan jalur aktual Anda. Contoh: /home/ssl/only.4096.client.truststore.jks.
PentingParameter ini tidak diperlukan jika Anda menggunakan titik akhir default atau titik akhir SASL. Parameter ini wajib jika Anda menggunakan titik akhir SSL.
ssl.truststore.password
Password truststore. Nilainya tetap KafkaOnsClient. Jangan ubah nilai ini.
ssl.endpoint.identification.algorithm
Algoritma verifikasi hostname. Biarkan kosong untuk menonaktifkan verifikasi hostname.
java.security.auth.login.config
Jalur file konfigurasi JAAS. Simpan file kafka_client_jaas.conf dari proyek demo ke direktori lokal dan ganti xxxx dengan jalur aktual. Contoh: /home/ssl/kafka_client_jaas.conf.
PentingParameter ini tidak diperlukan jika Anda menggunakan titik akhir default. Parameter ini wajib jika Anda menggunakan titik akhir SSL atau titik akhir SASL.
Kirim pesan
Kompilasi dan jalankan KafkaProducerDemo.java untuk mengirim pesan.
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.Future;
// Jika Anda menggunakan titik akhir SSL atau SASL, beri komentar pada baris berikut.
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
/*
* Jika Anda menggunakan titik akhir SSL atau SASL, hapus komentar pada dua baris berikut.
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
*/
public class KafkaProducerDemo {
public static void main(String args[]) {
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
* Atur jalur file konfigurasi JAAS.
JavaKafkaConfigurer.configureSasl();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada baris berikut.
* Atur jalur file konfigurasi JAAS.
JavaKafkaConfigurer.configureSaslPlain();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada baris berikut.
* Atur jalur file konfigurasi JAAS.
JavaKafkaConfigurer.configureSaslScram();
*/
// Muat kafka.properties.
Properties kafkaProperties = JavaKafkaConfigurer.getKafkaProperties();
Properties props = new Properties();
// Atur titik akhir. Peroleh titik akhir instans dari Konsol.
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getProperty("bootstrap.servers"));
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada empat baris berikut.
* Mirip dengan jalur SASL, file ini tidak dapat dikemas ke dalam file JAR.
props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, kafkaProperties.getProperty("ssl.truststore.location"));
* Password untuk truststore. Jangan ubah nilai ini.
props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
* Protokol keamanan. Untuk koneksi SSL, ini harus SASL_SSL.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
* Metode autentikasi SASL. Jangan ubah nilai ini.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada dua baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme PLAIN.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada dua baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme SCRAM.
props.put(SaslConfigs.SASL_MECHANISM, "SCRAM-SHA-256");
*/
// Serializer untuk kunci dan nilai pesan.
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
// Waktu tunggu maksimum untuk permintaan.
props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 30 * 1000);
// Jumlah retry di sisi client.
props.put(ProducerConfig.RETRIES_CONFIG, 5);
// Interval retry di sisi client.
props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 3000);
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
* Nilai kosong menonaktifkan verifikasi hostname.
props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");
*/
// Buat instance produsen. Produsen bersifat thread-safe. Satu instance produsen umumnya cukup untuk satu proses.
// Untuk meningkatkan performa, Anda dapat membuat lebih banyak instance produsen, tetapi disarankan tidak lebih dari lima.
KafkaProducer<String, String> producer = new KafkaProducer<String, String>(props);
// Buat pesan Kafka.
String topic = kafkaProperties.getProperty("topic"); // Topik untuk mengirim pesan. Anda harus membuat topik ini di Konsol terlebih dahulu.
String value = "this is the message's value"; // Konten pesan.
try {
// Mengirim pesan secara batch dan mengumpulkan future dapat meningkatkan performa, tetapi jangan membuat ukuran batch terlalu besar.
List<Future<RecordMetadata>> futures = new ArrayList<Future<RecordMetadata>>(128);
for (int i =0; i < 100; i++) {
// Kirim pesan dan peroleh objek Future.
ProducerRecord<String, String> kafkaMessage = new ProducerRecord<String, String>(topic, value + ": " + i);
Future<RecordMetadata> metadataFuture = producer.send(kafkaMessage);
futures.add(metadataFuture);
}
producer.flush();
for (Future<RecordMetadata> future: futures) {
// Peroleh hasil objek Future secara sinkron.
try {
RecordMetadata recordMetadata = future.get();
System.out.println("Produce ok:" + recordMetadata.toString());
} catch (Throwable t) {
t.printStackTrace();
}
}
} catch (Exception e) {
// Jika pesan gagal dikirim setelah retry di sisi client, aplikasi Anda harus menangani error ini.
System.out.println("error occurred");
e.printStackTrace();
}
}
}
Berlangganan pesan
Pilih salah satu metode berikut untuk berlangganan pesan.
Langganan konsumen tunggal
Kompilasi dan jalankan KafkaConsumerDemo.java untuk menerima pesan.
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.ProducerConfig;
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada tiga baris berikut. Jika Anda menggunakan titik akhir SASL, hapus komentar pada dua baris pertama.
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
*/
public class KafkaConsumerDemo {
public static void main(String args[]) {
// Atur jalur file konfigurasi JAAS.
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSasl();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSaslPlain();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSaslScram();
*/
// Muat kafka.properties.
Properties kafkaProperties = JavaKafkaConfigurer.getKafkaProperties();
Properties props = new Properties();
// Atur titik akhir. Peroleh titik akhir instans dari Konsol.
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getProperty("bootstrap.servers"));
// Jika Anda menggunakan titik akhir SSL, beri komentar pada baris berikut.
// Sesuaikan nilai ini berdasarkan jumlah data yang akan ditarik dan versi client. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada enam baris berikut.
* Jalur ke file truststore SSL. Pastikan jalur ini dikonfigurasi dengan benar di file kafka.properties Anda.
* Mirip dengan jalur SASL, file ini tidak dapat dikemas ke dalam file JAR.
props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, kafkaProperties.getProperty("ssl.truststore.location"));
* Password untuk truststore. Jangan ubah nilai ini.
props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
* Protokol keamanan. Untuk koneksi SSL, ini harus SASL_SSL.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
* Metode autentikasi SASL. Jangan ubah nilai ini.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
* Interval maksimum yang diizinkan antara dua polling.
* Jika konsumen gagal mengirim heartbeat dalam interval ini, broker menganggap konsumen tidak aktif, menghapusnya dari kelompok konsumen, dan memicu rebalance. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
* Jumlah data yang diambil per permintaan. Parameter ini dapat sangat memengaruhi performa saat mengakses instans melalui Internet.
props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 32000);
props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 32000);
*/
// Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, beri komentar pada baris berikut.
// Sesuaikan nilai ini berdasarkan jumlah data yang akan ditarik dan versi client. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada tiga baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme PLAIN.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
* Interval maksimum yang diizinkan antara dua polling.
* Jika konsumen gagal mengirim heartbeat dalam interval ini, broker menganggap konsumen tidak aktif, menghapusnya dari kelompok konsumen, dan memicu rebalance. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
*/
// Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, beri komentar pada baris berikut.
// Sesuaikan nilai ini berdasarkan jumlah data yang akan ditarik dan versi client. Nilai default: 30 detik
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada empat baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme SCRAM.
props.put(SaslConfigs.SASL_MECHANISM, "SCRAM-SHA-256");
* Interval maksimum yang diizinkan antara dua polling.
* Jika konsumen gagal mengirim heartbeat dalam interval ini, broker menganggap konsumen tidak aktif, menghapusnya dari kelompok konsumen, dan memicu rebalance. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
*/
// Jumlah maksimum catatan yang dikembalikan dalam satu panggilan poll().
// Jangan atur nilai ini terlalu besar. Jika Anda menarik data dalam jumlah besar tetapi tidak dapat mengonsumsi data tersebut sebelum polling berikutnya, rebalance akan dipicu, yang dapat menyebabkan lag.
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 30);
// Deserializer untuk kunci dan nilai pesan.
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
// Kelompok konsumen tempat instance konsumen saat ini berada. Anda harus membuat kelompok konsumen ini di Konsol terlebih dahulu.
// Instance konsumen dalam kelompok yang sama mengonsumsi pesan secara load-balanced.
props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaProperties.getProperty("group.id"));
// Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
// Nilai kosong menonaktifkan verifikasi hostname.
//props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");
// Buat instance konsumen.
KafkaConsumer<String, String> consumer = new org.apache.kafka.clients.consumer.KafkaConsumer<String, String>(props);
// Berlangganan topik. Anda dapat berlangganan beberapa topik.
// Jika instance konsumen memiliki nilai GROUP_ID_CONFIG yang sama, kami sarankan Anda mengonfigurasinya untuk berlangganan topik yang sama.
List<String> subscribedTopics = new ArrayList<String>();
// Jika Anda menggunakan titik akhir SSL, beri komentar pada lima baris pertama dan hapus komentar pada baris keenam.
// Jika Anda ingin berlangganan beberapa topik, tambahkan di sini.
// Anda harus membuat setiap topik di Konsol terlebih dahulu.
String topicStr = kafkaProperties.getProperty("topic");
String[] topics = topicStr.split(",");
for (String topic: topics) {
subscribedTopics.add(topic.trim());
}
//subscribedTopics.add(kafkaProperties.getProperty("topic"));
consumer.subscribe(subscribedTopics);
// Konsumsi pesan dalam loop.
while (true){
try {
ConsumerRecords<String, String> records = consumer.poll(1000);
// Data yang ditarik harus dikonsumsi sebelum polling berikutnya. Waktu total tidak boleh melebihi nilai SESSION_TIMEOUT_MS_CONFIG.
// Kami sarankan Anda membuat kolam thread terpisah untuk mengonsumsi pesan dan mengembalikan hasil secara asinkron.
for (ConsumerRecord<String, String> record : records) {
System.out.println(String.format("Consume partition:%d offset:%d", record.partition(), record.offset()));
}
} catch (Exception e) {
try {
Thread.sleep(1000);
} catch (Throwable ignore) {
}
e.printStackTrace();
}
}
}
}
Langganan multi-konsumen
Kompilasi dan jalankan KafkaMultiConsumerDemo.java untuk mengonsumsi pesan.
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicBoolean;
// Jika Anda menggunakan titik akhir SSL atau SASL, hapus komentar pada baris berikut.
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.ProducerConfig;
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada tiga baris berikut. Jika Anda menggunakan titik akhir SASL, hapus komentar pada dua baris pertama.
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
*/
import org.apache.kafka.common.errors.WakeupException;
/**
* Demo ini menunjukkan cara memulai beberapa konsumen dalam satu proses untuk mengonsumsi pesan dari suatu topik secara bersamaan.
* Pastikan jumlah total konsumen tidak melebihi jumlah total partisi untuk topik yang dilanggan.
*/
public class KafkaMultiConsumerDemo {
public static void main(String args[]) throws InterruptedException {
// Atur jalur file konfigurasi JAAS.
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSasl();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSaslPlain();
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada baris berikut.
JavaKafkaConfigurer.configureSaslScram();
*/
// Muat kafka.properties.
Properties kafkaProperties = JavaKafkaConfigurer.getKafkaProperties();
Properties props = new Properties();
// Atur titik akhir. Peroleh titik akhir instans dari Konsol.
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getProperty("bootstrap.servers"));
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada empat baris berikut.
* Mirip dengan jalur SASL, file ini tidak dapat dikemas ke dalam file JAR.
props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, kafkaProperties.getProperty("ssl.truststore.location"));
* Password untuk truststore. Jangan ubah nilai ini.
props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
* Protokol keamanan. Untuk koneksi SSL, ini harus SASL_SSL.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
* Metode autentikasi SASL. Jangan ubah nilai ini.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme PLAIN, hapus komentar pada dua baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme PLAIN.
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
*/
/*
* Jika Anda menggunakan titik akhir SASL dengan mekanisme SCRAM, hapus komentar pada dua baris berikut.
* Protokol keamanan.
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_PLAINTEXT");
* Mekanisme SCRAM.
props.put(SaslConfigs.SASL_MECHANISM, "SCRAM-SHA-256");
*/
// Interval maksimum yang diizinkan antara dua polling.
// Jika konsumen gagal mengirim heartbeat dalam interval ini, broker menganggap konsumen tidak aktif, menghapusnya dari kelompok konsumen, dan memicu rebalance. Nilai default: 30 detik.
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
// Jumlah maksimum catatan yang dikembalikan dalam satu panggilan poll().
// Jangan atur nilai ini terlalu besar. Jika Anda menarik data dalam jumlah besar tetapi tidak dapat mengonsumsi data tersebut sebelum polling berikutnya, rebalance akan dipicu, yang dapat menyebabkan lag.
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 30);
// Deserializer untuk kunci dan nilai pesan.
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
// Kelompok konsumen tempat instance konsumen saat ini berada. Anda harus membuat kelompok konsumen ini di Konsol terlebih dahulu.
// Instance konsumen dalam kelompok yang sama mengonsumsi pesan secara load-balanced.
props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaProperties.getProperty("group.id"));
/*
* Jika Anda menggunakan titik akhir SSL, hapus komentar pada baris berikut.
* Nilai kosong menonaktifkan verifikasi hostname.
props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");
*/
int consumerNum = 2;
Thread[] consumerThreads = new Thread[consumerNum];
for (int i = 0; i < consumerNum; i++) {
KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(props);
List<String> subscribedTopics = new ArrayList<String>();
subscribedTopics.add(kafkaProperties.getProperty("topic"));
consumer.subscribe(subscribedTopics);
KafkaConsumerRunner kafkaConsumerRunner = new KafkaConsumerRunner(consumer);
consumerThreads[i] = new Thread(kafkaConsumerRunner);
}
for (int i = 0; i < consumerNum; i++) {
consumerThreads[i].start();
}
for (int i = 0; i < consumerNum; i++) {
consumerThreads[i].join();
}
}
static class KafkaConsumerRunner implements Runnable {
private final AtomicBoolean closed = new AtomicBoolean(false);
private final KafkaConsumer consumer;
KafkaConsumerRunner(KafkaConsumer consumer) {
this.consumer = consumer;
}
@Override
public void run() {
try {
while (!closed.get()) {
try {
ConsumerRecords<String, String> records = consumer.poll(1000);
// Data yang ditarik harus dikonsumsi sebelum polling berikutnya. Waktu total tidak boleh melebihi nilai SESSION_TIMEOUT_MS_CONFIG.
for (ConsumerRecord<String, String> record : records) {
System.out.println(String.format("Thread:%s Consume partition:%d offset:%d", Thread.currentThread().getName(), record.partition(), record.offset()));
}
} catch (Exception e) {
try {
Thread.sleep(1000);
} catch (Throwable ignore) {
}
e.printStackTrace();
}
}
} catch (WakeupException e) {
// Abaikan exception jika konsumen ditutup.
if (!closed.get()) {
throw e;
}
} finally {
consumer.close();
}
}
// Hook shutdown yang dapat dipanggil oleh thread lain.
public void shutdown() {
closed.set(true);
consumer.wakeup();
}
}
}
FAQ
Konfigurasi sertifikat SASL_SSL
Unduh sertifikat SSL dari URL pada Langkah 1 bagian "Siapkan konfigurasi". Kemudian, di file kafka.properties Anda, atur ssl.truststore.location ke jalur lokal sertifikat tersebut.
Gunakan sertifikat SSL kustom
Tidak. Anda harus menggunakan sertifikat SSL yang disediakan oleh ApsaraMQ for Kafka.
Referensi
-
Anda juga dapat menggunakan framework Spring Cloud untuk menghubungkan ke instans. Untuk informasi selengkapnya, lihat Gunakan framework Spring Cloud untuk mengirim dan menerima pesan.
-
Jika Anda tidak dapat mengirim atau menerima pesan, periksa apakah instans dalam kondisi sehat. Untuk informasi selengkapnya, lihat Panduan pemeriksaan kesehatan instans.