All Products
Search
Document Center

DataHub:Buat subscription

Last Updated:Aug 25, 2026

Fitur subscription

Saat mengonsumsi data dari topik DataHub, Anda harus mengelola offset konsumsi secara manual agar dapat melanjutkan pemrosesan setelah terjadi kegagalan aplikasi. Hal ini mengharuskan Anda menyimpan progres konsumsi dan memastikan layanan penyimpanan offset memiliki ketersediaan tinggi, yang menambah kompleksitas pada aplikasi. Untuk menyederhanakan proses ini, DataHub menyediakan layanan subscription yang menyimpan offset konsumsi di sisi server. Dengan beberapa langkah konfigurasi sederhana dan kode minimal, Anda mendapatkan layanan manajemen offset berdaya tinggi yang beroperasi secara transparan terhadap aplikasi. Layanan subscription juga menyediakan kemampuan pengaturan ulang offset yang fleksibel, mendukung semantik konsumsi at-least-once. Misalnya, jika Anda menemukan kesalahan pemrosesan yang memengaruhi data dari periode waktu tertentu dan ingin mengonsumsi ulang data tersebut, Anda dapat mengatur ulang offset ke waktu yang sesuai. Aplikasi Anda secara otomatis mendeteksi perubahan ini dan memproses ulang data tanpa perlu melakukan restart.

Buat subscription

Pastikan akun Anda memiliki izin untuk membuat subscription pada topik dalam Proyek yang ditentukan. Untuk detail selengkapnya, lihat dokumentasi Pengendalian Izin. Ikuti langkah-langkah berikut:

  • Buka halaman Topic, klik + Subscription di pojok kanan atas, isi detail subscription, lalu klik Create.

    • Subscription Application: Nama aplikasi yang menggunakan subscription ini.

    • Description: Deskripsi lengkap tentang subscription.

  • Klik tombol pencarian di bawah Consumption Checkpoint untuk melihat status konsumsi semua shard.

Contoh penggunaan

Fitur subscription menyimpan offset. Meskipun fitur ini bersifat independen dari fungsi baca dan tulis DataHub (lihat dokumentasi Java SDK), fitur ini sering digunakan bersamaan dengannya ketika Anda perlu menyimpan offset konsumsi setelah membaca data.

// Contoh mengonsumsi data dan melakukan commit offset selama proses.
public void offset_consumption(int maxRetry) {
    String endpoint = "<YourEndPoint>";
    String accessId = "<YourAccessId>";
    String accessKey = "<YourAccessKey>";
    String projectName = "<YourProjectName>";
    String topicName = "<YourTopicName>";
    String subId = "<YourSubId>";
    String shardId = "0";
    List<String> shardIds = Arrays.asList(shardId);
    // Buat instans DatahubClient.
    DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
            .setDatahubConfig(
                    new DatahubConfig(endpoint,
                            // Apakah akan mengaktifkan transfer biner. Fitur ini didukung oleh server sejak versi 2.12.
                            new AliyunAccount(accessId, accessKey), true))
            .build();
    RecordSchema schema = datahubClient.getTopic(projectName, topicName).getRecordSchema();
    OpenSubscriptionSessionResult openSubscriptionSessionResult = datahubClient.openSubscriptionSession(projectName, topicName, subId, shardIds);
    SubscriptionOffset subscriptionOffset = openSubscriptionSessionResult.getOffsets().get(shardId);
    // 1. Dapatkan kursor untuk offset saat ini. Jika offset saat ini telah kedaluwarsa atau belum pernah dikonsumsi, dapatkan kursor untuk catatan pertama dalam siklus hidup.
    String cursor = "";
    // Nomor urut kurang dari 0 menunjukkan bahwa shard belum dikonsumsi.
    if (subscriptionOffset.getSequence() < 0) {
        // Dapatkan kursor untuk catatan pertama dalam siklus hidup.
        cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
    } else {
        // Dapatkan kursor untuk catatan berikutnya.
        long nextSequence = subscriptionOffset.getSequence() + 1;
        try {
            // Mendapatkan kursor menggunakan SEQUENCE dapat melempar SeekOutOfRangeException, yang menunjukkan bahwa data pada kursor saat ini telah kedaluwarsa.
            cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
        } catch (SeekOutOfRangeException e) {
            // Dapatkan kursor untuk catatan pertama dalam siklus hidup.
            cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
        }
    }
    // 2. Baca catatan dan simpan offset. Contoh ini menunjukkan pembacaan data tuple dan melakukan commit offset setiap 1.000 catatan.
    long recordCount = 0L;
    // Baca 1.000 catatan sekaligus.
    int fetchNum = 1000;
    int retryNum = 0;
    int commitNum = 1000;
    while (retryNum < maxRetry) {
        try {
            GetRecordsResult getRecordsResult = datahubClient.getRecords(projectName, topicName, shardId, schema, cursor, fetchNum);
            if (getRecordsResult.getRecordCount() <= 0) {
                // Tidak ada data. Tidur sejenak lalu coba lagi.
                System.out.println("no data, sleep 1 second");
                Thread.sleep(1000);
                continue;
            }
            for (RecordEntry recordEntry : getRecordsResult.getRecords()) {
                // Proses data.
                TupleRecordData data = (TupleRecordData) recordEntry.getRecordData();
                System.out.println("field1:" + data.getField("field1") + "\t"
                        + "field2:" + data.getField("field2"));
                // Setelah memproses data, perbarui offset.
                recordCount++;
                subscriptionOffset.setSequence(recordEntry.getSequence());
                subscriptionOffset.setTimestamp(recordEntry.getSystemTime());
                // Lakukan commit offset setiap 1000 catatan.
                if (recordCount % commitNum == 0) {
                    // Lakukan commit offset.
                    Map<String, SubscriptionOffset> offsetMap = new HashMap<>();
                    offsetMap.put(shardId, subscriptionOffset);
                    datahubClient.commitSubscriptionOffset(projectName, topicName, subId, offsetMap);
                    System.out.println("commit offset successful");
                }
            }
            cursor = getRecordsResult.getNextCursor();
        } catch (SubscriptionOfflineException | SubscriptionSessionInvalidException e) {
            // Keluar. SubscriptionOfflineException: Subscription sedang offline. SubscriptionSessionInvalidException: Klien lain sedang mengonsumsi subscription yang sama.
            e.printStackTrace();
            throw e;
        } catch (SubscriptionOffsetResetException e) {
            // Offset telah diatur ulang. Anda perlu mendapatkan versi terbaru dari objek SubscriptionOffset.
            SubscriptionOffset offset = datahubClient.getSubscriptionOffset(projectName, topicName, subId, shardIds).getOffsets().get(shardId);
            subscriptionOffset.setVersionId(offset.getVersionId());
            // Setelah offset diatur ulang, Anda harus mendapatkan kursor baru. Metode yang Anda gunakan untuk mendapatkan kursor harus sesuai dengan cara offset diatur ulang.
            // Jika sequence dan timestamp keduanya diatur saat pengaturan ulang, Anda dapat mendapatkan kursor menggunakan SEQUENCE atau SYSTEM_TIME.
            // Jika hanya sequence yang diatur, Anda harus menggunakan SEQUENCE.
            // Jika hanya timestamp yang diatur, Anda harus menggunakan SYSTEM_TIME.
            // Sebagai aturan umum, coba dapatkan kursor menggunakan SEQUENCE terlebih dahulu, lalu SYSTEM_TIME. Jika keduanya gagal, gunakan OLDEST.
            cursor = null;
            if (cursor == null) {
                try {
                    long nextSequence = offset.getSequence() + 1;
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by SEQUENCE failed, try to get cursor by SYSTEM_TIME");
                }
            }
            if (cursor == null) {
                try {
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SYSTEM_TIME, offset.getTimestamp()).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by SYSTEM_TIME failed, try to get cursor by OLDEST");
                }
            }
            if (cursor == null) {
                try {
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by OLDEST failed");
                    System.out.println("get cursor failed!!");
                    throw e;
                }
            }
        } catch (LimitExceededException e) {
            // Batas terlampaui, coba lagi.
            e.printStackTrace();
            retryNum++;
        } catch (DatahubClientException e) {
            // Kesalahan lain, coba lagi.
            e.printStackTrace();
            retryNum++;
        } catch (Exception e) {
            e.printStackTrace();
            System.exit(-1);
        }
    }
}
  • Saat aplikasi pertama kali dijalankan, aplikasi mulai mengonsumsi data dari catatan paling awal yang tersedia. Selama aplikasi berjalan, Anda dapat merefresh halaman subscription di Konsol Web untuk melihat offset konsumsi shard yang maju.

  • Jika Anda mengubah offset secara manual menggunakan fitur Reset Checkpoint di Konsol Web saat konsumen sedang berjalan, aplikasi secara otomatis mendeteksi perubahan tersebut dan melanjutkan konsumsi dari offset baru. Untuk melakukan ini, client menangkap SubscriptionOffsetResetException dan memanggil metode getSubscriptionOffset untuk mengambil objek SubscriptionOffset terbaru dari server.

  • Jangan gunakan beberapa thread atau proses konsumen untuk mengonsumsi shard yang sama dari satu subscription secara bersamaan. Hal ini menyebabkan offset ditimpa oleh konsumen yang berbeda, sehingga offset yang tersimpan berada dalam kondisi tidak terdefinisi. Dalam skenario ini, server melempar SubscriptionSessionInvalidException. Tangkap pengecualian ini, hentikan aplikasi, dan periksa desain Anda untuk menghindari konsumen duplikat.