Le SDK Tablestore pour Java permet de consommer en continu les données d'un tunnel, de traiter chaque lot via un rappel et de configurer les mécanismes de heartbeat, les points de contrôle (checkpoints), les pools de threads ainsi que le niveau de concurrence de la consommation.
Remarques d'utilisation
La période de rétention des journaux incrémentiels correspond à la durée d'expiration des journaux Stream de la table et peut atteindre sept jours. Pour un tunnel de type BaseAndStream, si la consommation des données complètes ne se termine pas dans ce délai, l'erreur
OTSTunnelExpiredest renvoyée au début de la consommation des données incrémentielles. Le tunnel ne peut alors plus continuer à consommer les données incrémentielles.Si la consommation incrémentielle prend du retard par rapport à la période de rétention, le tunnel peut reprendre à partir des dernières données disponibles. Par conséquent, certaines données risquent de ne pas être consommées.
Un tunnel expiré peut être désactivé. Si le tunnel reste désactivé pendant plus de 30 jours, il est supprimé et ne peut pas être restauré.
Prérequis
Installez le SDK Tablestore pour Java et initialisez TunnelClient.
Description de la fonctionnalité
L'objet TunnelWorker se connecte à un tunnel via son ID, utilise des heartbeats pour obtenir les canaux attribués au client actuel, extrait les données en continu et transmet chaque lot d'enregistrements à l'interface IChannelProcessor. Si plusieurs instances TunnelWorker consomment le même tunnel, le serveur répartit les canaux entre les clients.
Pour consommer des données à partir d'un tunnel :
Implémentez l'interface
IChannelProcessor. Utilisez la méthodeprocesspour traiter chaque lot et la méthodeshutdownpour libérer les ressources utilisées par le rappel.Créez un objet
TunnelWorkerConfigafin de configurer le rappel et le comportement de consommation.Créez une instance
TunnelWorkeren spécifiant l'ID du tunnel, l'objetTunnelClientet la configurationTunnelWorkerConfig.-
Appelez la méthode
connectAndWorkingpour démarrer la consommation.void process(ProcessRecordsInput input); void shutdown();
L'exemple suivant affiche chaque enregistrement extrait du tunnel, puis démarre la consommation.
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() {
// Release resources used by the callback.
}
}
String tunnelId = "example_tunnel_id";
TunnelWorkerConfig config =
new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker =
new TunnelWorker(tunnelId, tunnelClient, config);
worker.connectAndWorking();
La méthode connectAndWorking retourne après avoir lancé les tâches de consommation en arrière-plan. Maintenez le processus de l'application en cours d'exécution. Pour arrêter la consommation, appelez dans l'ordre les méthodes worker.shutdown(), config.shutdown() et tunnelClient.shutdown(). La méthode worker.shutdown() ferme les connexions au tunnel et invoque la méthode shutdown du rappel. La méthode config.shutdown() arrête les pools de threads dédiés à la lecture, au traitement et aux fonctions utilitaires. L'objet TunnelWorker enregistre un hook d'arrêt de la JVM qui tente d'arrêter le worker, mais l'application doit tout de même libérer explicitement ces ressources.
Paramètres
Worker
Le constructeur TunnelWorker accepte les paramètres suivants.
|
Nom |
Type |
Description |
|
tunnelId (obligatoire) |
String |
L'ID du tunnel. Obtenez-le en créant, en listant ou en interrogeant les tunnels. |
|
client (obligatoire) |
TunnelClientInterface |
L'instance |
|
workerConfig (obligatoire) |
TunnelWorkerConfig |
La configuration du rappel et du comportement de consommation. |
Configuration de la consommation
L'objet workerConfig est de type TunnelWorkerConfig et contient les paramètres suivants.
|
Nom |
Type |
Description |
|
channelProcessor (obligatoire) |
IChannelProcessor |
Le rappel de traitement des données. Ce paramètre est requis lorsque vous utilisez le constructeur |
|
heartbeatTimeoutInSec (facultatif) |
long |
Le délai d'expiration du heartbeat en secondes. La valeur par défaut est |
|
heartbeatIntervalInSec (facultatif) |
long |
L'intervalle entre les heartbeats en secondes. La valeur par défaut est |
|
checkpointIntervalInMillis (facultatif) |
long |
L'intervalle, en millisecondes, auquel les points de contrôle (checkpoints) de consommation sont enregistrés sur le serveur. La valeur par défaut est |
|
clientTag (facultatif) |
String |
Une étiquette client personnalisée utilisée pour générer un ID client et distinguer les instances |
|
readRecordsExecutor (facultatif) |
ThreadPoolExecutor |
Le pool de threads chargé d'extraire les données. Le pool par défaut compte |
|
processRecordsExecutor (facultatif) |
ThreadPoolExecutor |
Le pool de threads chargé de traiter les données. Sa configuration par défaut est identique à celle de |
|
maxChannelParallel (facultatif) |
int |
Le nombre maximal de canaux depuis lesquels les données sont extraites et traitées simultanément. Utilisez ce paramètre pour limiter l'utilisation de la mémoire. La valeur par défaut est |
|
channelHelperExecutor (facultatif) |
ThreadPoolExecutor |
Le pool de threads utilitaire qui initialise les canaux, planifie les pipelines et gère les erreurs d'exécution. Si ce paramètre n'est pas défini, un pool de threads mis en cache est utilisé. |
|
maxRetryIntervalInMillis (facultatif) |
int |
L'intervalle de base maximal pour la stratégie de backoff exponentiel lors de l'extraction des données incrémentielles, en millisecondes. La valeur par défaut est |
|
readMaxTimesPerRound (facultatif) |
int |
Le nombre maximal d'appels |
|
readMaxBytesPerRound (facultatif) |
int |
La quantité maximale de données extraites lors d'un cycle de pipeline, en octets. La valeur par défaut est |
|
enableClosingChannelDetect (facultatif) |
boolean |
Indique s'il faut détecter en temps réel les canaux dans l'état |
Si vous démarrez plusieurs instances TunnelWorker sur la même machine, vous pouvez réutiliser un seul objet TunnelWorkerConfig pour partager les pools de threads de lecture et de traitement. Une fois tous les workers arrêtés, appelez la méthode config.shutdown() une seule fois.
Données du rappel
La méthode process reçoit un objet ProcessRecordsInput qui contient les champs suivants.
|
Champ |
Type |
Description |
|
records |
|
Les enregistrements extraits dans le lot actuel. Appelez la méthode |
|
nextToken |
String |
Le jeton pour le lot suivant. Appelez la méthode |
|
traceId |
String |
L'ID de trace de la demande d'extraction actuelle. Appelez la méthode |
|
channelId |
String |
L'ID du canal auquel appartient le lot actuel. Appelez la méthode |
Exemples de scénarios
Ajuster les paramètres de consommation
Si le débit de consommation ou l'utilisation de la mémoire ne répond pas à vos exigences, ajustez les intervalles de heartbeat et de point de contrôle, la concurrence des canaux, le nombre et la taille des extractions par cycle, ainsi que l'intervalle de backoff pour les extractions incrémentielles.
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);