Tous les produits
Search
Centre de documentation

Tablestore:Consume data from a tunnel

Dernière mise à jour :Aug 20, 2026

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 OTSTunnelExpired est 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 :

  1. Implémentez l'interface IChannelProcessor. Utilisez la méthode process pour traiter chaque lot et la méthode shutdown pour libérer les ressources utilisées par le rappel.

  2. Créez un objet TunnelWorkerConfig afin de configurer le rappel et le comportement de consommation.

  3. Créez une instance TunnelWorker en spécifiant l'ID du tunnel, l'objet TunnelClient et la configuration TunnelWorkerConfig.

  4. Appelez la méthode connectAndWorking pour 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();
Important

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 TunnelClient initialisée.

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 TunnelWorker à trois paramètres.

heartbeatTimeoutInSec (facultatif)

long

Le délai d'expiration du heartbeat en secondes. La valeur par défaut est 300 et doit être supérieure à heartbeatIntervalInSec. Après l'expiration d'un heartbeat, le serveur considère le client comme indisponible et celui-ci se reconnecte au tunnel.

heartbeatIntervalInSec (facultatif)

long

L'intervalle entre les heartbeats en secondes. La valeur par défaut est 30 et la valeur minimale est 5. Les heartbeats permettent d'obtenir les canaux actifs, de mettre à jour l'état des canaux et d'initialiser les tâches de traitement des données. Cet intervalle influence également le temps de préchauffage de l'objet TunnelWorker.

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 5000. Le service Tunnel distribue chaque enregistrement au moins une fois et préserve l'ordre des enregistrements. Une tâche redémarrée reprend à partir du dernier point de contrôle, ce qui peut entraîner le traitement de certaines données plusieurs fois. Un intervalle plus court réduit les traitements en double, mais un enregistrement trop fréquent des points de contrôle peut diminuer le débit.

clientTag (facultatif)

String

Une étiquette client personnalisée utilisée pour générer un ID client et distinguer les instances TunnelWorker. La valeur par défaut correspond à la propriété système Java os.name.

readRecordsExecutor (facultatif)

ThreadPoolExecutor

Le pool de threads chargé d'extraire les données. Le pool par défaut compte 32 threads principaux, jusqu'à 1000 threads, une capacité de file d'attente de 16 et un temps de maintien en vie (keep-alive) de 60 secondes.

processRecordsExecutor (facultatif)

ThreadPoolExecutor

Le pool de threads chargé de traiter les données. Sa configuration par défaut est identique à celle de readRecordsExecutor. Pour un pool personnalisé, configurez le nombre de threads en fonction du nombre de canaux du tunnel.

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 -1, ce qui signifie qu'il n'y a aucune limite. Le SDK Tablestore pour Java version 5.10.0 et ultérieures prend en charge ce paramètre.

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 2000 et la valeur minimale est 200. Si un lot contient au maximum 500 enregistrements et ne dépasse pas 900 Ko, le client augmente progressivement l'intervalle de backoff. L'intervalle réel est sélectionné aléatoirement entre 75 % et 125 % de l'intervalle de base actuel. Le SDK Tablestore pour Java version 5.4.0 et ultérieures prend en charge ce paramètre.

readMaxTimesPerRound (facultatif)

int

Le nombre maximal d'appels ReadRecords par cycle de pipeline. La valeur par défaut est 1.

readMaxBytesPerRound (facultatif)

int

La quantité maximale de données extraites lors d'un cycle de pipeline, en octets. La valeur par défaut est 4194304, soit 4 MiB. Le cycle s'arrête lorsque cette valeur ou readMaxTimesPerRound est atteinte.

enableClosingChannelDetect (facultatif)

boolean

Indique s'il faut détecter en temps réel les canaux dans l'état CLOSING. Un canal CLOSING est en cours de migration d'un client vers un autre. Le SDK Tablestore pour Java version 5.13.13 et ultérieures prend en charge ce paramètre. La valeur par défaut est true dans la version 5.17.0 et ultérieures. Si la détection est désactivée, la migration des canaux peut être bloquée et la consommation interrompue lorsqu'il existe de nombreux canaux mais que les ressources client sont insuffisantes.

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

List<StreamRecord>

Les enregistrements extraits dans le lot actuel. Appelez la méthode getRecords() pour les obtenir.

nextToken

String

Le jeton pour le lot suivant. Appelez la méthode getNextToken() pour l'obtenir. L'objet TunnelWorker utilise automatiquement cette valeur pour poursuivre l'extraction des données et enregistrer les points de contrôle.

traceId

String

L'ID de trace de la demande d'extraction actuelle. Appelez la méthode getTraceId() pour l'obtenir.

channelId

String

L'ID du canal auquel appartient le lot actuel. Appelez la méthode getChannelId() pour l'obtenir. Appelez la méthode getPartitionId() pour obtenir l'ID de partition à partir de l'ID de canal.

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);