Todos os produtos
Search
Central de documentação

Tablestore:Consume data from a tunnel

Última atualização: Aug 20, 2026

O Tablestore SDK for Java consome dados continuamente de um túnel, processa cada lote com um callback e permite configurar heartbeats, checkpoints, thread pools e concorrência de consumo.

Observações de uso

  • O período de retenção de logs incrementais corresponde ao período de expiração do log de Stream da tabela e pode chegar a sete dias. Em túneis do tipo BaseAndStream, se o consumo completo dos dados não terminar dentro desse prazo, o sistema retornará o erro OTSTunnelExpired ao iniciar o consumo de dados incrementais. Nesse caso, o túnel não poderá continuar consumindo dados incrementais.

  • Se o consumo incremental ficar atrasado em relação ao período de retenção, o túnel poderá retomar a partir dos dados disponíveis mais recentes. Consequentemente, alguns dados podem não ser consumidos.

  • Um túnel expirado pode ser desativado. Se permanecer desativado por mais de 30 dias, ele será excluído e não poderá ser restaurado.

Pré-requisitos

Instale o Tablestore SDK for Java e inicialize o TunnelClient.

Descrição do recurso

O TunnelWorker conecta-se a um túnel pelo ID do túnel, usa heartbeats para obter os canais atribuídos ao cliente atual, puxa dados continuamente e repassa cada lote de registros para o IChannelProcessor. Quando várias instâncias do TunnelWorker consomem o mesmo túnel, o servidor distribui os canais entre os clientes.

Para consumir dados de um túnel:

  1. Implemente o IChannelProcessor. Use o método process para processar cada lote e o método shutdown para liberar os recursos usados pelo callback.

  2. Crie um objeto TunnelWorkerConfig para configurar o callback e o comportamento de consumo.

  3. Crie um TunnelWorker informando o ID do túnel, o TunnelClient e o TunnelWorkerConfig.

  4. Chame o método connectAndWorking para iniciar o consumo.

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

O exemplo a seguir imprime cada registro obtido do túnel e inicia o consumo.

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

O método connectAndWorking retorna assim que as tarefas de consumo em segundo plano são iniciadas. Mantenha o processo da aplicação em execução. Para interromper o consumo, chame worker.shutdown(), config.shutdown() e tunnelClient.shutdown() nessa ordem. O método worker.shutdown() fecha as conexões do túnel e invoca o método shutdown do callback. O config.shutdown() encerra os thread pools de leitura, processamento e auxiliares. Embora o TunnelWorker registre um hook de desligamento da JVM que tenta parar o worker, a aplicação ainda precisa liberar esses recursos explicitamente.

Parâmetros

Worker

O construtor do TunnelWorker aceita os seguintes parâmetros.

Nome

Tipo

Descrição

tunnelId (obrigatório)

String

ID do túnel. Obtido ao criar, listar ou consultar túneis.

client (obrigatório)

TunnelClientInterface

Instância inicializada do TunnelClient.

workerConfig (obrigatório)

TunnelWorkerConfig

Configuração do callback e do comportamento de consumo.

Configuração de consumo

O parâmetro workerConfig é do tipo TunnelWorkerConfig e contém os seguintes parâmetros.

Nome

Tipo

Descrição

channelProcessor (obrigatório)

IChannelProcessor

Callback para processamento de dados. Obrigatório ao usar o construtor de três parâmetros do TunnelWorker.

heartbeatTimeoutInSec (opcional)

long

Tempo limite do heartbeat em segundos. O valor padrão é 300 e deve ser maior que heartbeatIntervalInSec. Após o tempo limite do heartbeat, o servidor considera o cliente indisponível e o cliente se reconecta ao túnel.

heartbeatIntervalInSec (opcional)

long

Intervalo do heartbeat em segundos. O valor padrão é 30 e o mínimo é 5. Os heartbeats obtêm canais ativos, atualizam estados dos canais e inicializam tarefas de processamento de dados. Esse intervalo também afeta o tempo de aquecimento do TunnelWorker.

checkpointIntervalInMillis (opcional)

long

Intervalo em milissegundos para gravação de checkpoints de consumo no servidor. O valor padrão é 5000. O Tunnel service entrega cada registro pelo menos uma vez e preserva a ordem dos registros. Uma tarefa reiniciada retoma a partir do último checkpoint e pode processar alguns dados mais de uma vez. Um intervalo menor reduz o processamento duplicado, mas registrar checkpoints com muita frequência pode diminuir o throughput.

clientTag (opcional)

String

Tag personalizada do cliente usada para gerar um ID de cliente e distinguir instâncias do TunnelWorker. O valor padrão é a propriedade de sistema Java os.name.

readRecordsExecutor (opcional)

ThreadPoolExecutor

Thread pool responsável por puxar dados. O pool padrão possui 32 threads principais, até 1000 threads, capacidade de fila de 16 e tempo de keep-alive de 60 segundos.

processRecordsExecutor (opcional)

ThreadPoolExecutor

Thread pool destinado ao processamento de dados. Sua configuração padrão é idêntica à do readRecordsExecutor. Para um pool personalizado, configure o número de threads com base na quantidade de canais do túnel.

maxChannelParallel (opcional)

int

Número máximo de canais dos quais os dados são puxados e processados simultaneamente. Use este parâmetro para limitar o uso de memória. O valor padrão é -1, que indica ausência de limite. O Tablestore SDK for Java 5.10.0 e versões posteriores suportam este parâmetro.

channelHelperExecutor (opcional)

ThreadPoolExecutor

Thread pool auxiliar que inicializa canais, agenda pipelines e trata erros de runtime. Se este parâmetro não for definido, um thread pool em cache será utilizado.

maxRetryIntervalInMillis (opcional)

int

Intervalo base máximo para backoff exponencial durante a extração de dados incrementais, em milissegundos. O valor padrão é 2000 e o mínimo é 200. Se um lote contiver no máximo 500 registros e não ultrapassar 900 KB, o cliente aumenta gradualmente o intervalo de backoff. O intervalo real é selecionado aleatoriamente entre 75% e 125% do intervalo base atual. O Tablestore SDK for Java 5.4.0 e versões posteriores suportam este parâmetro.

readMaxTimesPerRound (opcional)

int

Quantidade máxima de chamadas de ReadRecords em uma única rodada do pipeline. O valor padrão é 1.

readMaxBytesPerRound (opcional)

int

Volume máximo de dados extraídos em uma rodada do pipeline, em bytes. O valor padrão é 4194304, ou 4 MiB. A rodada é interrompida quando esse valor ou o readMaxTimesPerRound é atingido.

enableClosingChannelDetect (opcional)

boolean

Defina se deve detectar canais no estado CLOSING em tempo real. Um canal em CLOSING está sendo migrado de um cliente para outro. O Tablestore SDK for Java 5.13.13 e versões posteriores suportam este parâmetro. O valor padrão é true na versão 5.17.0 e posteriores. Se a detecção estiver desativada, a migração de canais pode ser bloqueada e o consumo interrompido quando houver muitos canais e os recursos do cliente forem insuficientes.

Ao iniciar várias instâncias do TunnelWorker na mesma máquina, reutilize um único objeto TunnelWorkerConfig para compartilhar os thread pools de leitura e processamento. Depois que todos os workers pararem, chame config.shutdown() apenas uma vez.

Dados do callback

O método process recebe um objeto ProcessRecordsInput com os seguintes campos.

Campo

Tipo

Descrição

records

List<StreamRecord>

Registros extraídos no lote atual. Chame getRecords() para obtê-los.

nextToken

String

Token referente ao próximo lote. Chame getNextToken() para obtê-lo. O TunnelWorker usa esse valor automaticamente para continuar puxando dados e registrar checkpoints.

traceId

String

ID de rastreamento da solicitação de extração atual. Chame getTraceId() para obtê-lo.

channelId

String

ID do canal ao qual o lote atual pertence. Chame getChannelId() para obtê-lo. Chame getPartitionId() para obter o ID da partição a partir do ID do canal.

Exemplos de cenários

Ajustar parâmetros de consumo

Se o throughput de consumo ou o uso de memória não atender aos seus requisitos, ajuste os intervalos de heartbeat e checkpoint, a concorrência de canais, a quantidade e o tamanho das extrações por rodada e o intervalo de backoff para extrações incrementais.

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