Após criar uma tarefa de rastreamento de alterações, use o kit de desenvolvimento de software (SDK) fornecido pelo Data Transmission Service (DTS) para assinar as mudanças nos dados. Este tópico descreve como usar um SDK para consumir dados de fontes distribuídas, como PolarDB-X 1.0 e bancos de dados lógicos do DMS.
Pré-requisitos
Crie uma instância de rastreamento de alterações com estado Normal. Para mais informações, consulte Create a change tracking task for a PolarDB-X 1.0 instance ou Create a change tracking task for a DMS logical database.
Crie um created a consumer group para sua instância de assinatura.
Se você utilizar um usuário RAM para consumir os dados assinados, conceda a ele a permissão AliyunDTSFullAccess e acesso aos objetos assinados. Para mais detalhes, consulte Grant permissions to a RAM user to manage DTS using a system policy e Manage the permissions of a RAM user.
Observações de uso
Ao consumir dados assinados, chame o método commit de DefaultUserRecord para confirmar as informações de offset. Caso contrário, pode ocorrer consumo duplicado dos dados.
Os processos de consumo são independentes entre si.
Procedimento
Baixe e descompacte o código de exemplo do SDK.
-
Verifique a versão do código do SDK.
Acesse o diretório onde você descompactou o código de exemplo do SDK.
Abra o arquivo pom.xml nesse diretório com um editor de texto.
-
Atualize o SDK de rastreamento de alterações para a versão mais recente.
NotaEncontre a dependência Maven mais recente na página do dts-new-subscribe-sdk.
-
Edite o código do SDK.
Abra o arquivo descompactado com um editor de código.
-
Conforme o padrão de uso do cliente SDK, abra o arquivo DistributedDTSConsumerDemo.java.
NotaO caminho do arquivo Java é
aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/. -
Defina os parâmetros no código Java.
public static void main(String[] args) throws ClientException { // Configuration for subscribing to a distributed data source, such as PolarDB-X 1.0 (formerly DRDS). Configure information such as the AccessKey, instance ID, main task ID, and consumer group. 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 seek(a timestamp to set, eg 1566180200 if you want (Mon Aug 19 10:03:21 CST 2019)) String checkpoint = "1639620090"; // Convert physical database/table name to logical database/table name boolean mapping = true; // if force use config checkpoint when start. for checkpoint reset, only assign mode works 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(); }Parâmetro
Descrição
Como obter
accessKeyId
O AccessKey ID.
Para mais informações, consulte Obtain an AccessKey pair.
accessKeySecret
O AccessKey secret.
regionId
O ID da região onde a tarefa de rastreamento de alterações está localizada.
No console do DTS, clique em no ID da instância de rastreamento de alterações desejada. Na página Basic Information, obtenha as informações da região. Por exemplo, se a região for China (Hangzhou), defina este parâmetro como
cn-hangzhou. Para mais detalhes, consulte List of regions.dtsInstanceId
O ID da instância de rastreamento de alterações.
No console do DTS, clique em no ID da instância de rastreamento de alterações desejada. Na página Basic Information, obtenha o DTS Instance ID da instância.
jobId
O ID da tarefa de rastreamento de alterações.
Chame a operação DescribeDtsJobs para obter o ID da tarefa de rastreamento de alterações (DtsJobId).
sid
O ID do grupo de consumidores.
No console do DTS, clique em no ID da instância de rastreamento de alterações desejada. No painel de navegação à esquerda, clique em Consume Data. Obtenha o Consumer Group ID/Name e a Account do grupo de consumidores.
NotaA senha da conta do grupo de consumidores é definida durante a criação do grupo.
userName
A conta do grupo de consumidores.
password
A senha da conta do grupo de consumidores.
proxyUrl
O endpoint e a porta do canal de rastreamento de alterações.
NotaSe a instância ECS onde o cliente SDK está implantado e o canal de rastreamento estiverem na mesma rede clássica ou Virtual Private Cloud (VPC), assine os dados pela rede interna para obter a menor latência possível.
-
Não recomendamos o uso de endpoint público devido à possível instabilidade da rede.
No console do DTS, clique em no ID da instância de rastreamento de alterações desejada. Na página Basic Information, obtenha as informações de Network.
checkpoint
O offset de consumo. Trata-se do timestamp a partir do qual o cliente SDK começa a consumir registros de dados. O valor é um timestamp UNIX em segundos.
NotaAs informações de offset de consumo podem ser usadas para:
Retomar o consumo de dados e evitar perda de informações caso o processo seja interrompido, bastando passar o offset de consumo.
Ajustar o offset de assinatura ao iniciar o cliente SDK, passando o offset desejado para consumir dados conforme necessário.
O offset de consumo deve estar dentro do intervalo de timestamps da instância de rastreamento de alterações e convertido para timestamp UNIX.
NotaVisualize o intervalo de timestamps da instância na coluna Data Range da lista de tarefas de rastreamento.
Utilize um mecanismo de busca para encontrar um conversor de timestamp UNIX.
-
Opcional: Para modificar o tipo de dado dos dados assinados, altere o método
buildRecordListener()ou utilize uma classe personalizada.public static Map<String, RecordListener> buildRecordListener() { // user can impl their own 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)) { // consume record RecordListener recordPrintListener = new DefaultRecordPrintListener(DbType.MySQL); recordPrintListener.consume(record); //commit method push the checkpoint update record.commit(""); } } }; return Collections.singletonMap("mysqlRecordPrinter", mysqlRecordPrintListener); } Abra a estrutura do projeto na sua IDE e verifique se a versão do OpenJDK do projeto é a 1,8.
-
Execute o código do cliente.
A saída indica que o cliente está assinando alterações de dados do banco de dados de origem.
-
Periodicamente, o cliente SDK coleta e exibe estatísticas sobre o consumo de dados. Essas estatísticas incluem o número total de registros enviados e recebidos, o volume total de dados e os registros por segundo (RPS).
Tabela 1. Estatísticas de consumo de dados
Parâmetro
Descrição
outCountsNúmero total de registros de dados consumidos pelo cliente SDK.
outBytesVolume total de dados consumidos pelo cliente SDK, em bytes.
outRpsQuantidade de solicitações por segundo enviadas pelo cliente SDK para consumir dados.
outBpsTaxa de bits transmitidos por segundo durante o consumo de dados pelo cliente SDK.
countNenhum.
inBytesVolume total de dados enviados pelo servidor DTS, em bytes.
DStoreRecordQueueTamanho da fila de cache de dados quando o servidor DTS envia dados.
inCountsNúmero total de registros de dados enviados pelo servidor DTS.
inRpsQuantidade de solicitações enviadas pelo servidor DTS por segundo.
inBpsTaxa de bits transmitidos por segundo quando o servidor DTS envia dados.
__dtTimestamp em que o cliente SDK recebe os dados, em milissegundos.
DefaultUserRecordQueueTamanho da fila de cache de dados após a serialização.