Após criar uma tarefa de rastreamento de alterações, use a demonstração de SDK fornecida pelo Data Transmission Service (DTS) para consumir as mudanças de dados resultantes. Este tópico explica como usar a demonstração de SDK com fontes de dados distribuídas, como PolarDB-X 1.0 e DMS LogicDB.
Pré-requisitos
JDK 1.8 instalado.
IntelliJ IDEA instalado.
Precauções
Para consumir dados de alteração como usuário RAM, conceda a esse usuário a permissão AliyunDTSFullAccess e permissões para acessar os objetos de source. Para mais informações, consulte Autorizar um usuário RAM a gerenciar o DTS e Gerencie permissões de usuário RAM.
Procedimento
Este tópico demonstra como execute a demonstração de SDK para consumir dados de alteração de uma instância PolarDB-X 1.0, usando o IntelliJ IDEA Community Edition 2020.1 para Windows como exemplo.
Crie uma instância de rastreamento de alterações. Para mais informações, consulte Criar uma tarefa de rastreamento de alterações para PolarDB-X 1.0.
Crie um ou mais grupos de consumidores. Para mais informações, consulte Criar um grupo de consumidores.
-
Baixe e descompacte a demonstração de SDK. Para obter o link de download, consulte Código da Demonstração de SDK.
ImportanteAo consumir dados de alteração, chame o método commit de DefaultUserRecord para envie o offset do consumidor. Caso contrário, há risco de consumir dados repetidamente.
-
Abra o projeto no IntelliJ IDEA.
Abra o IntelliJ IDEA e clique em Open or Import.
Na caixa de diálogo, acesse o diretório onde a demonstração de SDK foi descompactada, expanda as pastas e clique duas vezes no arquivo pom.xml.
Na caixa de diálogo exibida, selecione Open as Project.
-
No IntelliJ IDEA, expanda as pastas. Conforme o modo de uso do cliente SDK, selecione e clique duas vezes para abrir o arquivo Java correspondente: DistributedDTSConsumerDemo.
aliyun-dts-subscribe-sdk-java-master [dts-new-subscr...] .idea src main test java com.aliyun.dts.subscribe.clients DBMapperTest DistributedDTSConsumerDemo DTSConsumerAssignDemo DTSConsumerSubscribeDemo UserMetaStore .gitignore dts-new-subscribe-sdk.iml LICENSE pom.xml README.md External Libraries Scratches and Consoles -
Defina os parâmetros necessários no arquivo Java.
public static void main(String[] args) throws ClientException { // Configure the change tracking settings for a distributed data source, such as a PolarDB-X 1.0 instance. Set parameters such as the AccessKey pair, instance ID, 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"; // The initial consumer offset, a Unix timestamp in seconds. e.g., 1566180200 for Mon Aug 19 10:03:21 CST 2019. String checkpoint = "1639620090"; // Convert physical database/table name to logical database/table name boolean mapping = true; // Set to true to force the client to start from the specified checkpoint. A checkpoint reset works only 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(); }Parâmetro
Descrição
Como obter
accessKeyId
O AccessKey id.
Para mais informações, consulte Obter um par de AccessKey.
accessKeySecret
O AccessKey secret.
regionId
O id da região da instância de rastreamento de alterações.
No console DTS, clique em no id da instância de rastreamento de alterações desejada. Na página Task Management, localize as informações da região. Por exemplo, se a região for China (Hangzhou), defina este parâmetro como
cn-hangzhou. Para obter uma lista de regiões, consulte Lista de regiões suportadas.dtsInstanceId
O id da instância de rastreamento de alterações.
No console DTS, clique em no id da instância de rastreamento de alterações desejada. Na página Task Management, localize o id da instância e o id da tarefa.
jobId
O id da tarefa de rastreamento de alterações.
sid
O id do grupo de consumidores.
No console DTS, clique em no id da instância de rastreamento de alterações desejada e clique em Consume Data. Localize o SID 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 da instância de rastreamento de alterações.
Nota-
Para minimizar a latência de rede, use o endpoint interno se a instância ECS que executa o cliente SDK estiver na mesma rede clássica ou VPC da instância de rastreamento de alterações.
-
Não recomendamos o uso de endpoint público devido à possível instabilidade da rede.
No console DTS, clique em no id da instância de rastreamento de alterações desejada. Na página Task Management, localize o endpoint e o número da porta.
checkpoint
O offset do consumidor, especificado como timestamp Unix em segundos. O cliente SDK começa a consumir dados a partir deste momento.
NotaO offset do consumidor é usado nos seguintes cenários:
-
Se a aplicação for interrompida, passe o último offset do consumidor para retomar o consumo sem perda de dados.
-
Ao iniciar o cliente, passe um offset específico para consumir dados sob demanda.
O offset do consumidor deve estar dentro do intervalo de dados da instância de rastreamento de alterações e ser um timestamp Unix.
NotaUse um mecanismo de busca para encontrar um conversor de timestamp Unix.
-
-
Na barra de menu superior do IntelliJ IDEA, escolha para execute o cliente.
NotaNa primeira execução, o IntelliJ IDEA instala automaticamente as dependências necessárias, o que pode levar algum tempo.
A saída indica que o cliente consegue consumir alterações de dados da instância de source.
-
O cliente SDK exibe periodicamente estatísticas de consumo, incluindo o volume total e a quantidade de registros de dados enviados e recebidos, além das requisições por segundo (RPS).
Tabela 1. Estatísticas de consumo
Parâmetro
Descrição
outCountsQuantidade total de registros de dados consumidos pelo cliente SDK.
outBytesVolume total de dados consumidos pelo cliente SDK, em bytes.
outRpsTaxa de consumo de dados do cliente SDK, em requisições por segundo (RPS).
outBpsTaxa de consumo de dados do cliente SDK, em bits por segundo (bps).
countParâmetro reservado.
inBytesVolume total de dados enviados pelo servidor DTS, em bytes.
DStoreRecordQueueTamanho atual da fila de cache de dados no servidor DTS.
inCountsQuantidade total de registros de dados enviados pelo servidor DTS.
inRpsTaxa de envio de dados do servidor DTS, em requisições por segundo (RPS).
inBpsTaxa de envio de dados do servidor DTS, em bits por segundo (bps).
__dtTimestamp em que o cliente SDK recebe os dados, em milissegundos.
DefaultUserRecordQueueTamanho atual da fila de cache de dados após a serialização.
-
Opcional: Para personalizar o processamento dos dados consumidos, modifique o método
buildRecordListener()ou use uma classe personalizada.public static Map<String, RecordListener> buildRecordListener() { // You can implement your 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 the record. RecordListener recordPrintListener = new DefaultRecordPrintListener(DbType.MySQL); recordPrintListener.consume(record); // The commit method pushes the checkpoint update. record.commit(""); } } }; return Collections.singletonMap("mysqlRecordPrinter", mysqlRecordPrintListener); }