All Products
Search
Document Center

DataHub:Kompatibilitas Kafka

Last Updated:Aug 27, 2026

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_topic dipetakan ke proyektest_project dan 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 project.topic:subId, ID tersebut harus sesuai dengan topik yang dilanggan. Jika tidak, data tidak dapat dibaca. Kami menyarankan menggunakan format project.group.

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

Tidak

Strategi penetapan partisi default di Kafka adalah RangeAssignor. DataHub saat ini hanya mendukung strategi ini. Jangan ubah parameter ini.

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. Karenasession.timeout.ms disesuaikan menjadi 60.000 ms, kamimenyarankan agar Anda secara eksplisit mengatur nilai ini menjadi 40000 untuk mencegah permintaan heartbeat yang terlalu sering.

application.id

S

project.topic:subId

atau

project.group

Ya

Jika Anda menggunakan format project.topic:subId, ID tersebut harus sesuai dengan topik yang dilanggan. Jika tidak, data tidak dapat dibaca. Kami menyarankan menggunakan format project.group.

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

  1. Buat topik di Konsol

    Saat membuat topik, aktifkan Shard Expand Mode.

  2. Buat topik menggunakan SDK

    Anda tidak dapat membuat topik menggunakan API Kafka. Anda harus menggunakan SDK DataHub dan mengatur ExpandMode ke ONLY_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

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

  2. 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) sebagai group.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 found

Tambahkan 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.DisconnectException

Klien 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.