All Products
Search
Document Center

Tablestore:Mengonsumsi data dari saluran data

Last Updated:Aug 01, 2026

Tablestore SDK for Java dapat terus-menerus mengonsumsi data dari saluran data, memproses setiap batch melalui panggilan balik (callback), serta mengonfigurasi heartbeat, checkpoint, kolam thread, dan konkurensi konsumsi.

Catatan penggunaan

  • Periode retensi log inkremental sama dengan periode kedaluwarsa log Stream tabel dan dapat mencapai tujuh hari. Untuk saluran data BaseAndStream, jika konsumsi data penuh tidak selesai dalam periode tersebut, error OTSTunnelExpired dikembalikan saat konsumsi data inkremental dimulai. Saluran data tersebut tidak dapat melanjutkan konsumsi data inkremental.

  • Jika konsumsi inkremental tertinggal di luar periode retensi, saluran data mungkin melanjutkan dari data terbaru yang tersedia. Akibatnya, beberapa data mungkin tidak terkonsumsi.

  • Saluran data yang kedaluwarsa mungkin dinonaktifkan. Jika tetap dinonaktifkan lebih dari 30 hari, saluran tersebut akan dihapus dan tidak dapat dipulihkan.

Prasyarat

Instal Tablestore SDK for Java dan inisialisasi TunnelClient.

Deskripsi fitur

TunnelWorker terhubung ke saluran data berdasarkan ID saluran data, menggunakan heartbeat untuk mendapatkan saluran yang ditugaskan ke klien saat ini, terus-menerus menarik data, dan meneruskan setiap batch catatan ke IChannelProcessor. Jika beberapa instans TunnelWorker mengonsumsi saluran data yang sama, server mendistribusikan saluran tersebut di antara klien-klien tersebut.

Untuk mengonsumsi data dari saluran data:

  1. Implementasikan IChannelProcessor. Gunakan metode process untuk memproses setiap batch dan metode shutdown untuk melepaskan sumber daya yang digunakan oleh callback.

  2. Buat objek TunnelWorkerConfig untuk mengonfigurasi callback dan perilaku konsumsi.

  3. Buat TunnelWorker dengan ID saluran data, TunnelClient, dan TunnelWorkerConfig.

  4. Panggil connectAndWorking untuk memulai konsumsi.

    void process(ProcessRecordsInput input);
    void shutdown();

Contoh berikut mencetak setiap catatan yang ditarik dari saluran data, lalu memulai konsumsi.

private static class SimpleProcessor implements IChannelProcessor {
    @Override
    public void process(ProcessRecordsInput input) {
        for (StreamRecord record : input.getRecords()) {
            System.out.println(record);
        }
    }

    @Override
    public void shutdown() {
        // Melepaskan sumber daya yang digunakan oleh callback.
    }
}

String tunnelId = "example_tunnel_id";
TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker =
        new TunnelWorker(tunnelId, tunnelClient, config);
worker.connectAndWorking();
Penting

connectAndWorking mengembalikan nilai setelah tugas konsumsi latar belakang dimulai. Pastikan proses aplikasi tetap berjalan. Untuk menghentikan konsumsi, panggil worker.shutdown(), config.shutdown(), dan tunnelClient.shutdown() secara berurutan. worker.shutdown() menutup koneksi saluran data dan memanggil metode shutdown pada callback. config.shutdown() menghentikan kolam thread pembacaan, pemrosesan, dan helper. Meskipun TunnelWorker mendaftarkan hook shutdown JVM yang mencoba menghentikan worker, aplikasi tetap harus secara eksplisit melepaskan sumber daya tersebut.

Parameter

Worker

Konstruktor TunnelWorker memiliki parameter berikut.

Nama

Type

Deskripsi

tunnelId (wajib)

String

ID saluran data. Dapatkan melalui pembuatan, pencatatan, atau kueri saluran data.

client (wajib)

TunnelClientInterface

TunnelClient yang telah diinisialisasi.

workerConfig (wajib)

TunnelWorkerConfig

Konfigurasi callback dan perilaku konsumsi.

Konfigurasi konsumsi

workerConfig bertipe TunnelWorkerConfig dan memiliki parameter berikut.

Nama

Type

Deskripsi

channelProcessor (wajib)

IChannelProcessor

Callback pemrosesan data. Parameter ini wajib saat menggunakan konstruktor TunnelWorker dengan tiga parameter.

heartbeatTimeoutInSec (opsional)

long

Timeout heartbeat dalam detik. Nilai default adalah 300 dan harus lebih besar dari heartbeatIntervalInSec. Setelah heartbeat timeout, server menganggap klien tidak tersedia dan klien akan terhubung ulang ke saluran data.

heartbeatIntervalInSec (opsional)

long

Interval heartbeat dalam detik. Nilai default adalah 30, dan nilai minimum adalah 5. Heartbeat digunakan untuk mendapatkan saluran aktif, memperbarui status saluran, dan menginisialisasi tugas pemrosesan data. Interval ini juga memengaruhi waktu pemanasan (warm-up time) TunnelWorker.

checkpointIntervalInMillis (opsional)

long

Interval pencatatan titik pemeriksaan konsumsi di server, dalam milidetik. Nilai default adalah 5000. Tunnel Service mengirimkan setiap catatan minimal sekali dan mempertahankan urutan catatan. Tugas yang dimulai ulang dilanjutkan dari checkpoint terbaru, sehingga beberapa data mungkin diproses lebih dari sekali. Interval yang lebih pendek mengurangi pemrosesan duplikat, tetapi terlalu sering mencatat checkpoint dapat mengurangi throughput.

clientTag (opsional)

String

Tag klien kustom yang digunakan untuk menghasilkan ID klien dan membedakan instans TunnelWorker. Nilai default adalah properti sistem Java os.name.

readRecordsExecutor (opsional)

ThreadPoolExecutor

Kolam thread yang menarik data. Kolam default memiliki 32 thread inti, hingga maksimal 1000 thread, kapasitas antrian 16, dan waktu hidup thread (keep-alive time) 60 detik.

processRecordsExecutor (opsional)

ThreadPoolExecutor

Kolam thread yang memproses data. Konfigurasi default-nya sama dengan readRecordsExecutor. Untuk kolam kustom, atur jumlah thread berdasarkan jumlah saluran dalam saluran data.

maxChannelParallel (opsional)

int

Jumlah maksimum saluran tempat data ditarik dan diproses secara konkuren. Gunakan parameter ini untuk membatasi penggunaan memori. Nilai default adalah -1, yang berarti tanpa batas. Parameter ini didukung mulai Tablestore SDK for Java versi 5.10.0.

channelHelperExecutor (opsional)

ThreadPoolExecutor

Kolam thread helper yang menginisialisasi saluran, menjadwalkan pipeline, dan menangani error waktu proses. Jika parameter ini tidak diatur, kolam thread cache digunakan.

maxRetryIntervalInMillis (opsional)

int

Interval dasar maksimum untuk backoff eksponensial selama penarikan data inkremental, dalam milidetik. Nilai default adalah 2000, dan nilai minimum adalah 200. Jika suatu batch berisi tidak lebih dari 500 catatan dan ukurannya tidak melebihi 900 KB, klien secara bertahap meningkatkan interval backoff. Interval aktual dipilih secara acak antara 75% hingga 125% dari interval dasar saat ini. Parameter ini didukung mulai Tablestore SDK for Java versi 5.4.0.

readMaxTimesPerRound (opsional)

int

Jumlah maksimum panggilan ReadRecords dalam satu putaran pipeline. Nilai default adalah 1.

readMaxBytesPerRound (opsional)

int

Jumlah maksimum data yang ditarik dalam satu putaran pipeline, dalam byte. Nilai default adalah 4194304, atau 4 MiB. Putaran berhenti ketika nilai ini atau readMaxTimesPerRound tercapai.

enableClosingChannelDetect (opsional)

boolean

Menentukan apakah akan mendeteksi saluran dalam status CLOSING secara real-time. Saluran CLOSING sedang dalam proses migrasi dari satu klien ke klien lain. Parameter ini didukung mulai Tablestore SDK for Java versi 5.13.13. Nilai default adalah true pada versi 5.17.0 dan seterusnya. Jika deteksi dinonaktifkan, migrasi saluran dapat terblokir dan konsumsi dapat terganggu ketika banyak saluran ada tetapi sumber daya klien tidak mencukupi.

Jika Anda menjalankan beberapa instans TunnelWorker pada mesin yang sama, Anda dapat menggunakan kembali satu objek TunnelWorkerConfig untuk berbagi kolam thread pembacaan dan pemrosesan. Setelah semua worker berhenti, panggil config.shutdown() hanya sekali.

Data callback

Metode process menerima objek ProcessRecordsInput yang berisi bidang-bidang berikut.

Bidang

Type

Deskripsi

records

List<StreamRecord>

Catatan yang ditarik dalam batch saat ini. Panggil getRecords() untuk mendapatkannya.

nextToken

String

Token untuk batch berikutnya. Panggil getNextToken() untuk mendapatkannya. TunnelWorker secara otomatis menggunakan nilai ini untuk melanjutkan penarikan data dan mencatat checkpoint.

traceId

String

ID jejak dari permintaan tarik saat ini. Panggil getTraceId() untuk mendapatkannya.

channelId

String

ID saluran tempat batch saat ini berasal. Panggil getChannelId() untuk mendapatkannya. Panggil getPartitionId() untuk mendapatkan ID partisi dari ID saluran.

Contoh skenario

Menyesuaikan parameter konsumsi

Jika throughput konsumsi atau penggunaan memori tidak memenuhi kebutuhan Anda, sesuaikan interval heartbeat dan checkpoint, konkurensi saluran, jumlah dan ukuran penarikan per putaran, serta interval backoff untuk penarikan inkremental.

TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
config.setHeartbeatIntervalInSec(10);
config.setHeartbeatTimeoutInSec(60);
config.setCheckpointIntervalInMillis(10_000);
config.setMaxChannelParallel(16);
config.setReadMaxTimesPerRound(4);
config.setReadMaxBytesPerRound(8 * 1024 * 1024);
config.setMaxRetryIntervalInMillis(3_000);