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
OTSTunnelExpireddikembalikan 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:
-
Implementasikan
IChannelProcessor. Gunakan metodeprocessuntuk memproses setiap batch dan metodeshutdownuntuk melepaskan sumber daya yang digunakan oleh callback. -
Buat objek
TunnelWorkerConfiguntuk mengonfigurasi callback dan perilaku konsumsi. -
Buat
TunnelWorkerdengan ID saluran data,TunnelClient, danTunnelWorkerConfig. -
Panggil
connectAndWorkinguntuk 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();
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 |
|
|
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 |
|
heartbeatTimeoutInSec (opsional) |
long |
Timeout heartbeat dalam detik. Nilai default adalah |
|
heartbeatIntervalInSec (opsional) |
long |
Interval heartbeat dalam detik. Nilai default adalah |
|
checkpointIntervalInMillis (opsional) |
long |
Interval pencatatan titik pemeriksaan konsumsi di server, dalam milidetik. Nilai default adalah |
|
clientTag (opsional) |
String |
Tag klien kustom yang digunakan untuk menghasilkan ID klien dan membedakan instans |
|
readRecordsExecutor (opsional) |
ThreadPoolExecutor |
Kolam thread yang menarik data. Kolam default memiliki |
|
processRecordsExecutor (opsional) |
ThreadPoolExecutor |
Kolam thread yang memproses data. Konfigurasi default-nya sama dengan |
|
maxChannelParallel (opsional) |
int |
Jumlah maksimum saluran tempat data ditarik dan diproses secara konkuren. Gunakan parameter ini untuk membatasi penggunaan memori. Nilai default adalah |
|
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 |
|
readMaxTimesPerRound (opsional) |
int |
Jumlah maksimum panggilan |
|
readMaxBytesPerRound (opsional) |
int |
Jumlah maksimum data yang ditarik dalam satu putaran pipeline, dalam byte. Nilai default adalah |
|
enableClosingChannelDetect (opsional) |
boolean |
Menentukan apakah akan mendeteksi saluran dalam status |
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 |
|
Catatan yang ditarik dalam batch saat ini. Panggil |
|
nextToken |
String |
Token untuk batch berikutnya. Panggil |
|
traceId |
String |
ID jejak dari permintaan tarik saat ini. Panggil |
|
channelId |
String |
ID saluran tempat batch saat ini berasal. Panggil |
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);