O Data Transmission Service (DTS) oferece um demo do sdk Java para consumo de dados com rastreamento de alterações em bancos de dados distribuídos. Este guia orienta você a baixe, configure e execute o demo do sdk em uma instância PolarDB for Xscale (PolarDB-X) 1.0 ou em um banco de dados lógico do Data Management (DMS).
Pré-requisitos
Antes de começar, verifique se você possui:
Java Development Kit (JDK) V1.8 instale
IntelliJ IDEA instale
Uma tarefa de rastreamento de alterações configure no DTS — consulte Rastrear alterações de dados de uma instância ApsaraDB RDS for MySQL
Um ou mais grupos de consumidores crie — consulte Criar grupos de consumidores
Permissões de usuário RAM
Para rastrear e consumir dados como usuário do Resource Access Management (RAM), é necessário que o usuário tenha a permissão AliyunDTSFullAccess e permissões para acessar os objetos de source. Para obter detalhes, consulte Usar uma política de sistema para autorizar um usuário RAM a gerencie instâncias do DTS e Conceder permissões ao usuário RAM.
Configure e execute o demo do sdk
Os passos a seguir utilizam o IntelliJ IDEA Community Edition 2020.1 para Windows.
Passo 1: Baixe o demo do sdk
Baixe o pacote do demo do sdk e descompacte-o.
Passo 2: Abrir o projeto no IntelliJ IDEA
-
Abra o IntelliJ IDEA. Na janela de boas-vindas, clique em Open or Import.

-
Acesse o diretório onde você descompactou o pacote, abra as pastas e clique duas vezes em pom.xml.

Na caixa de diálogo exibida, selecione Open as Project.
Passo 3: Abrir DistributedDTSConsumerDemo
No IntelliJ IDEA, expanda as pastas do projeto para localizar os arquivos Java e clique duas vezes em DistributedDTSConsumerDemo.

Passo 4: Configure os parâmetros necessários
Defina os parâmetros no método main() da classe DistributedDTSConsumerDemo:
public static void main(String[] args) throws ClientException {
// Configure a change tracking task for a distributed database such as a PolarDB-X 1.0 instance.
// Set the parameters for your AccessKey pair, instance ID, task ID, and consumer groups.
String accessKeyId = "LTA***********99reZ";
String accessKeySecret = "****************";
String regionId = "cn-hangzhou";
String dtsInstanceId = "dtse5212sed162****";
String jobId = "l791216x16d****";
String sid = "dtsip412t13160****";
String userName = "xftest";
String password = "******";
String proxyUrl = "dts-cn-****.com:18001";
// Initial checkpoint for first consumption (UNIX timestamp in seconds).
// Example: 1566180200 corresponds to Mon Aug 19 10:03:21 CST 2019.
String checkpoint = "1639620090";
// Convert physical database/table names to logical database/table names.
boolean mapping = true;
// Force-use the initial checkpoint when starting. Only works in ASSIGN mode.
boolean isForceUseInitCheckpoint = false;
ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.ASSIGN;
DistributedDTSConsumerDemo demo = new DistributedDTSConsumerDemo(userName, password, regionId,
jobId, sid, dtsInstanceId, accessKeyId, accessKeySecret, subscribeMode, proxyUrl,
checkpoint, isForceUseInitCheckpoint, mapping);
demo.start();
}
A tabela a seguir descreve cada parâmetro.
|
Parâmetro |
Descrição |
Como obter |
|
|
accesskey ID |
Consulte Criar e obter um par de AccessKeys. |
|
|
AccessKey secret |
Consulte Criar e obter um par de AccessKeys. |
|
|
ID da região da tarefa de rastreamento de alterações. Exemplo: |
No console do DTS, clique em no ID da instância de rastreamento de alterações. Na página Basic Information, localize a região da instância. Para todos os IDs de região válidos, consulte Regiões suportadas. |
|
|
ID da instância de rastreamento de alterações. |
No console do DTS, clique em no ID da instância. O DTS Instance ID está listado na página Basic Information. |
|
|
ID da tarefa de rastreamento de alterações. |
Chame a operação da API DescribeDtsJobs para consultar o ID da tarefa. |
|
|
ID do grupo de consumidores. |
No console do DTS, clique em no ID da instância. No painel de navegação à esquerda, clique em Consume Data. Localize o Consumer Group ID/Name na lista. |
|
|
Conta do grupo de consumidores. |
No console do DTS, clique em no ID da instância. No painel de navegação à esquerda, clique em Consume Data. Encontre a Account na lista. |
|
|
Senha da conta do grupo de consumidores. |
A senha é defina durante a criação do grupo de consumidores. |
|
|
Endpoint e porta da instância de rastreamento de alterações. |
No console do DTS, clique em no ID da instância. Na página Basic Information, localize Network. Para reduzir a latência, implante o cliente do sdk em uma instância elastic compute service (ecs) que utilize a rede clássica ou compartilhe a mesma Virtual Private Cloud (VPC) da instância de rastreamento de alterações. |
|
|
Offset do consumidor — um timestamp UNIX em segundos que controla onde o cliente do sdk inicia o consumo de dados. Utilize este parâmetro para retomar após uma interrupção e evitar perda de dados, ou para iniciar o consumo em um timestamp personalizado conforme suas necessidades de negócio. |
Encontre o intervalo de dados da instância de rastreamento de alterações na coluna Data Range da página Change Tracking Tasks. O checkpoint deve estar dentro desse intervalo. Use um conversor de timestamp UNIX para obter o valor exato. |
Passo 5: Execute o cliente
Na barra de menu superior do IntelliJ IDEA, escolha Run > Run.
Na primeira execução do cliente, o IntelliJ IDEA pode levar alguns minutos para baixe e instale as dependências do Maven.
Quando o cliente for iniciado com sucesso, o console mostrará que as alterações de dados estão sendo rastreadas da instância de source. O cliente também relata métricas de consumo em intervalos regulares. A tabela a seguir descreve essas métricas.
|
Métrica |
Descrição |
|
|
Número total de registros de dados consumidos pelo cliente do sdk |
|
|
Quantidade total de dados consumidos pelo cliente do sdk, em bytes |
|
|
Taxa de consumo do cliente do sdk, em requisições por segundo (RPS) |
|
|
Throughput de consumo do cliente do sdk, em bits por segundo |
|
|
Nenhum |
|
|
Quantidade total de dados enviados pelo servidor DTS, em bytes |
|
|
Tamanho da fila de cache de dados no lado do servidor DTS |
|
|
Número total de registros de dados enviados pelo servidor DTS |
|
|
Taxa de envio do servidor DTS, em RPS |
|
|
Throughput de envio do servidor DTS, em bits por segundo |
|
|
Timestamp em que o cliente do sdk recebeu os dados, em milissegundos |
|
|
Tamanho da fila de cache de dados após a serialização |
Personalizar o listener de registros (opcional)
Por padrão, o demo do sdk imprime os registros consumidos no console. Para processar os registros de forma diferente — por exemplo, gravando-os em um banco de dados ou encaminhando-os para uma fila de mensagens — modifique o método buildRecordListener() ou implemente uma classe RecordListener personalizada.
O sdk entrega quatro tipos de operação: INSERT, UPDATE, DELETE e HEARTBEAT. O exemplo a seguir mostra como lidar com elas:
public static Map<String, RecordListener> buildRecordListener() {
// Implement your own record listener.
RecordListener mysqlRecordPrintListener = new RecordListener() {
@Override
public void consume(DefaultUserRecord record) {
OperationType operationType = record.getOperationType();
if (operationType.equals(OperationType.INSERT)
|| operationType.equals(OperationType.UPDATE)
|| operationType.equals(OperationType.DELETE)
|| operationType.equals(OperationType.HEARTBEAT)) {
// Process the record.
RecordListener recordPrintListener = new DefaultRecordPrintListener(DbType.MySQL);
recordPrintListener.consume(record);
// Commit pushes the checkpoint update to DTS.
record.commit("");
}
}
};
return Collections.singletonMap("mysqlRecordPrinter", mysqlRecordPrintListener);
}