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
OTSTunnelExpiredao 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:
Implemente o
IChannelProcessor. Use o métodoprocesspara processar cada lote e o métodoshutdownpara liberar os recursos usados pelo callback.Crie um objeto
TunnelWorkerConfigpara configurar o callback e o comportamento de consumo.Crie um
TunnelWorkerinformando o ID do túnel, oTunnelCliente oTunnelWorkerConfig.-
Chame o método
connectAndWorkingpara 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();
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 |
|
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 |
|
heartbeatTimeoutInSec (opcional) |
long |
Tempo limite do heartbeat em segundos. O valor padrão é |
|
heartbeatIntervalInSec (opcional) |
long |
Intervalo do heartbeat em segundos. O valor padrão é |
|
checkpointIntervalInMillis (opcional) |
long |
Intervalo em milissegundos para gravação de checkpoints de consumo no servidor. O valor padrão é |
|
clientTag (opcional) |
String |
Tag personalizada do cliente usada para gerar um ID de cliente e distinguir instâncias do |
|
readRecordsExecutor (opcional) |
ThreadPoolExecutor |
Thread pool responsável por puxar dados. O pool padrão possui |
|
processRecordsExecutor (opcional) |
ThreadPoolExecutor |
Thread pool destinado ao processamento de dados. Sua configuração padrão é idêntica à do |
|
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 é |
|
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 é |
|
readMaxTimesPerRound (opcional) |
int |
Quantidade máxima de chamadas de |
|
readMaxBytesPerRound (opcional) |
int |
Volume máximo de dados extraídos em uma rodada do pipeline, em bytes. O valor padrão é |
|
enableClosingChannelDetect (opcional) |
boolean |
Defina se deve detectar canais no estado |
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 |
|
Registros extraídos no lote atual. Chame |
|
nextToken |
String |
Token referente ao próximo lote. Chame |
|
traceId |
String |
ID de rastreamento da solicitação de extração atual. Chame |
|
channelId |
String |
ID do canal ao qual o lote atual pertence. Chame |
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);