O Tunnel Service permite consumir dados de uma tabela. Este tópico descreve como começar a usar o Tunnel Service com o Tablestore SDK for Java.
Observações de uso
Por padrão, o sistema inicia um pool de threads para ler e processar dados com base no TunnelWorkerConfig. Para iniciar vários TunnelWorkers em um único servidor, use o mesmo TunnelWorkerConfig na configuração de todos eles.
O TunnelWorker exige um período de aquecimento para inicialização, definido pelo parâmetro heartbeatIntervalInSec no TunnelWorkerConfig. Use o método setHeartbeatIntervalInSec no TunnelWorkerConfig para definir esse parâmetro. Valor padrão: 30. Unidade: segundos.
Se o cliente TunnelWorker for encerrado devido a uma saída inesperada ou término manual, ele reciclará automaticamente os recursos por meio de um dos seguintes métodos: liberação do pool de threads, chamada automática do método shutdown registrado na classe Channel e encerramento do túnel.
O período de retenção de logs incrementais nos túneis corresponde ao período de retenção de logs do Stream. Como os logs do Stream podem ser retidos por até sete dias, os logs incrementais nos túneis também têm retenção máxima de sete dias.
-
Ao criar um túnel para consumir dados diferenciais ou incrementais, observe os seguintes pontos:
-
Durante o consumo completo de dados, se o túnel não concluir o consumo dentro do período de retenção dos logs incrementais (máximo de sete dias), ocorrerá um erro
OTSTunnelExpiredao iniciar o consumo de logs incrementais. Consequentemente, o túnel não conseguirá consumir esses logs.Caso estime que o túnel não concluirá o consumo completo de dados dentro da janela de tempo especificada, entre em contato com o suporte técnico do Tablestore .
Durante o consumo de dados incrementais, se o túnel não concluir o consumo dos logs incrementais dentro do período de retenção (máximo de sete dias), ele poderá consumir dados a partir dos registros disponíveis mais recentes. Nesse cenário, alguns dados específicos podem não ser consumidos.
-
Após a expiração de um túnel, o Tablestore pode desativá-lo. Se permanecer no estado desativado por mais de 30 dias, o túnel será excluído. Não é possível restaurar um túnel excluído.
Pré-requisitos
-
Realize as seguintes operações no console Resource Access Management (RAM):
-
Crie um usuário RAM e anexe a política
AliyunOTSFullAccessa ele para conceder permissões de gerenciamento do Tablestore. Para obter mais informações, consulte Criar um usuário RAM e Gerenciar permissões de usuário RAM.NotaEm ambientes de produção, recomendamos conceder apenas as permissões necessárias aos usuários RAM, seguindo o princípio do menor privilégio. Essa prática ajuda a prevenir riscos de segurança causados por permissões excessivas.
-
Crie um par de AccessKey para o usuário RAM. Para obter mais informações, consulte Criar um par de AccessKey.
AvisoSe o par de AccessKey da sua conta Alibaba Cloud vazar, todos os recursos da conta ficarão expostos a riscos potenciais. Recomendamos usar o par de AccessKey de um usuário RAM para realizar operações, evitando assim o vazamento do par de AccessKey da sua conta Alibaba Cloud.
-
-
Realize as seguintes operações no console do Tablestore:
Crie uma tabela de dados. Para obter mais informações, consulte Usar o console do Tablestore, Usar o Tablestore CLI e Usar SDKs do Tablestore.
Obtenha o endpoint da instância onde a tabela de dados reside. Para obter mais informações, consulte Inicializar um cliente Tablestore.
Configure as credenciais de acesso. Para obter mais informações, consulte Configurar credenciais de acesso.
Começar a usar o Tunnel Service
Use o Tablestore SDK for Java para começar a utilizar o Tunnel Service.
-
Inicialize um cliente de túnel.
NotaCertifique-se de que as variáveis de ambiente
TABLESTORE_ACCESS_KEY_IDeTABLESTORE_ACCESS_KEY_SECRETestejam configuradas. A variável TABLESTORE_ACCESS_KEY_ID especifica o AccessKey ID da sua conta Alibaba Cloud ou usuário RAM. A variável TABLESTORE_ACCESS_KEY_SECRET especifica o AccessKey secret da sua conta Alibaba Cloud ou usuário RAM.// Specify the name of the Tablestore instance. // Specify the endpoint of the Tablestore instance. Example: https://instance.cn-hangzhou.ots.aliyuncs.com. // Specify the AccessKey ID and AccessKey secret of your Alibaba Cloud account or RAM user. final String instanceName = "yourInstanceName"; final String endPoint = "yourEndpoint"; final String accessKeyId = System.getenv("TABLESTORE_ACCESS_KEY_ID"); final String accessKeySecret = System.getenv("TABLESTORE_ACCESS_KEY_SECRET"); TunnelClient tunnelClient = new TunnelClient(endPoint, accessKeyId, accessKeySecret, instanceName); -
Crie um túnel.
Antes de criar um túnel, crie uma tabela de dados para teste ou prepare uma tabela existente. Crie a tabela no console do Tablestore ou usando o método createTable de um SyncClient.
ImportanteAo criar um túnel do tipo Stream ou BaseAndStream, siga as regras abaixo para especificar os timestamps:
Se você não especificar o timestamp inicial para os dados incrementais, o timestamp inicial será o momento da criação do túnel.
-
Se você especificar o timestamp inicial e o timestamp final para os dados incrementais, o timestamp inicial ou final deve estar dentro do intervalo [Hora atual do sistema - Período de validade do Stream + 5 minutos, Hora atual do sistema] em milissegundos.
O período de validade do Stream refere-se à validade dos logs incrementais em milissegundos. O período máximo é de sete dias. Defina esse período ao ativar o Stream para a tabela de dados. Após a definição, não é possível modificar o período de validade do Stream.
O timestamp final deve ser maior que o timestamp inicial.
// You can create three types of tunnels: TunnelType.BaseData, TunnelType.Stream, and TunnelType.BaseAndStream. // The following code provides an example on how to create a BaseAndStream tunnel. To create a tunnel of a different type, set TunnelType in CreateTunnelRequest to the desired type. final String tableName = "testTable"; final String tunnelName = "testTunnel"; CreateTunnelRequest request = new CreateTunnelRequest(tableName, tunnelName, TunnelType.BaseAndStream); CreateTunnelResponse resp = tunnelClient.createTunnel(request); // Use tunnelId to initialize TunnelWorker. Call the ListTunnel or DescribeTunnel operation to obtain the tunnel ID. String tunnelId = resp.getTunnelId(); System.out.println("Create Tunnel, Id: " + tunnelId); -
Especifique um callback personalizado de consumo de dados para iniciar o consumo automático. A tabela a seguir descreve as configurações do TunnelClient.
// Specify the callback for data consumption to call the IChannelProcessor operation, which specifies the process and shutdown methods. private static class SimpleProcessor implements IChannelProcessor { @Override public void process(ProcessRecordsInput input) { // The ProcessRecordsInput parameter contains the data that you obtain. System.out.println("Default record processor, would print records count"); System.out.println( // The NextToken parameter is used by the tunnel client to paginate data. String.format("Process %d records, NextToken: %s", input.getRecords().size(), input.getNextToken())); try { // Simulate data consumption and processing. Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); } } @Override public void shutdown() { System.out.println("Mock shutdown"); } } // By default, the system starts a thread pool to read and process data based on TunnelWorkerConfig. If you use a single server, multiple TunnelWorkers are started. // We recommend that you configure TunnelWorkers by using the same TunnelWorkerConfig. TunnelWorkerConfig provides more advanced parameters. TunnelWorkerConfig config = new TunnelWorkerConfig(new SimpleProcessor()); // Configure TunnelWorkers and start automatic data processing. TunnelWorker worker = new TunnelWorker(tunnelId, tunnelClient, config); try { worker.connectAndWorking(); } catch (Exception e) { e.printStackTrace(); config.shutdown(); worker.shutdown(); tunnelClient.shutdown(); }
Configurar o TunnelWorkerConfig
O TunnelWorkerConfig permite definir parâmetros personalizados para um cliente de túnel conforme suas necessidades de negócio. A tabela a seguir descreve os parâmetros disponíveis no Tablestore SDK for Java.
|
Configuração |
Parâmetro |
Descrição |
|
Intervalo e tempo limite para heartbeats |
heartbeatTimeoutInSec |
Tempo limite para heartbeats. Valor padrão: 300. Unidade: segundos. Quando ocorre um tempo limite de heartbeat, o servidor de túnel considera a instância atual do TunnelClient indisponível. Nesse caso, o cliente de túnel precisa se reconectar ao servidor. |
|
heartbeatIntervalInSec |
Intervalo entre heartbeats. Valor padrão: 30. Valor mínimo: 5. Unidade: segundos. A detecção de heartbeats monitora canais ativos, atualiza o status dos canais e inicializa automaticamente tarefas de processamento de dados. |
|
|
Intervalo entre checkpoints |
checkpointIntervalInMillis |
Intervalo entre checkpoints durante o consumo de dados. O intervalo é registrado no servidor de túnel. Valor padrão: 5000. Unidade: milissegundos. Nota
|
|
Tag personalizada do cliente |
clientTag |
Tag personalizada do cliente usada para gerar um ID de cliente de túnel. Configure este parâmetro para distinguir diferentes TunnelWorkers. |
|
Callbacks personalizados para processamento de dados |
channelProcessor |
Callback registrado pelo usuário para processar dados, incluindo os métodos process e shutdown. |
|
Configuração dos pools de threads para leitura e processamento de dados |
readRecordsExecutor |
Pool de threads usado para ler dados. Se não houver requisitos especiais, utilize a configuração padrão. |
|
processRecordsExecutor |
Pool de threads usado para processar dados. Se não houver requisitos especiais, utilize a configuração padrão. Nota
|
|
|
Controle de memória |
maxChannelParallel |
Nível máximo de concorrência de canais para leitura e processamento de dados, visando o controle de memória. O valor padrão é -1, indicando que a concorrência é ilimitada. Nota
O Tablestore SDK for Java V5.10.0 e versões posteriores suportam este recurso. |
|
Tempo máximo de backoff |
maxRetryIntervalInMillis |
Valor base para calcular o tempo máximo de backoff do túnel. O tempo máximo de backoff é um número aleatório entre 0,75 × maxRetryIntervalInMillis e 1,25 × maxRetryIntervalInMillis. Valor padrão: 2000. Valor mínimo: 200. Unidade: milissegundos. Nota
|
|
Detecção de canais CLOSING |
enableClosingChannelDetect |
Define se a detecção em tempo real para canais CLOSING deve ser ativada. Valor padrão: false, indicando que a detecção em tempo real para canais CLOSING está desativada. Nota
|
Apêndice: Código de exemplo
import com.alicloud.openservices.tablestore.TunnelClient;
import com.alicloud.openservices.tablestore.model.tunnel.CreateTunnelRequest;
import com.alicloud.openservices.tablestore.model.tunnel.CreateTunnelResponse;
import com.alicloud.openservices.tablestore.model.tunnel.TunnelType;
import com.alicloud.openservices.tablestore.tunnel.worker.IChannelProcessor;
import com.alicloud.openservices.tablestore.tunnel.worker.ProcessRecordsInput;
import com.alicloud.openservices.tablestore.tunnel.worker.TunnelWorker;
import com.alicloud.openservices.tablestore.tunnel.worker.TunnelWorkerConfig;
public class TunnelQuickStart {
private static class SimpleProcessor implements IChannelProcessor {
@Override
public void process(ProcessRecordsInput input) {
System.out.println("Default record processor, would print records count");
System.out.println(
// The NextToken parameter is used to by the Tunnel client to paginate data.
String.format("Process %d records, NextToken: %s", input.getRecords().size(), input.getNextToken()));
try {
// Simulate data consumption and processing.
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
@Override
public void shutdown() {
System.out.println("Mock shutdown");
}
}
public static void main(String[] args) throws Exception {
//1. Initialize the tunnel client.
// Specify the name of the instance.
final String instanceName = "yourInstanceName";
// Specify the endpoint of the instance.
final String endPoint = "yourEndpoint";
// Obtain the AccessKey ID and AccessKey secret from the environment variables.
final String accessKeyId = System.getenv("TABLESTORE_ACCESS_KEY_ID");
final String accessKeySecret = System.getenv("TABLESTORE_ACCESS_KEY_SECRET");
TunnelClient tunnelClient = new TunnelClient(endPoint, accessKeyId, accessKeySecret, instanceName);
//2. Create a tunnel. Before you perform this step, you must create a table for testing. You can create a table in the Tablestore console or by using the createTable method of a SyncClient.
final String tableName = "testTable";
final String tunnelName = "testTunnel";
CreateTunnelRequest request = new CreateTunnelRequest(tableName, tunnelName, TunnelType.BaseAndStream);
CreateTunnelResponse resp = tunnelClient.createTunnel(request);
// Use the tunnelId parameter to initialize a TunnelWorker. You can call the ListTunnel or DescribeTunnel operation to obtain the tunnel ID.
String tunnelId = resp.getTunnelId();
System.out.println("Create Tunnel, Id: " + tunnelId);
//3. Specify a custom callback function to start automatic data consumption.
// TunnelWorkerConfig provides more advanced parameters.
TunnelWorkerConfig config = new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker = new TunnelWorker(tunnelId, tunnelClient, config);
try {
worker.connectAndWorking();
} catch (Exception e) {
e.printStackTrace();
config.shutdown();
worker.shutdown();
tunnelClient.shutdown();
}
}
}