Todos os produtos
Search
Central de documentação

Data Transmission Service:Usar o SDK para consumir dados rastreados de uma instância PolarDB-X 1.0

Última atualização: Aug 21, 2026

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

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

  1. Baixe e descompacte o código de exemplo do SDK.

  2. Verifique a versão do código do SDK.

    1. Acesse o diretório onde você descompactou o código de exemplo do SDK.

    2. Abra o arquivo pom.xml nesse diretório com um editor de texto.

    3. Atualize o SDK de rastreamento de alterações para a versão mais recente.

      Nota

      Encontre a dependência Maven mais recente na página do dts-new-subscribe-sdk.

      Localização do parâmetro de versão do SDK (clique em para expandir)

      <name>dts-new-subscribe-sdk</name>
      <url>https://www.aliyun.com/product/dts</url>
      <description>The Aliyun new Subscribe SDK for Java used for accessing Data Transmission Service</description>
      <packaging>jar</packaging>
      <groupId>com.aliyun.dts</groupId>
      <artifactId>dts-new-subscribe-sdk</artifactId>
      <version>2.1.4</version>
  3. Edite o código do SDK.

    1. Abra o arquivo descompactado com um editor de código.

    2. Conforme o padrão de uso do cliente SDK, abra o arquivo DistributedDTSConsumerDemo.java.

      Nota

      O caminho do arquivo Java é aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/.

    3. 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.

      Nota

      A 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.

      Nota
      • Se 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.

      Nota

      As 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.

      Nota
      • Visualize 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.

  4. 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);
        }
  5. Abra a estrutura do projeto na sua IDE e verifique se a versão do OpenJDK do projeto é a 1,8.

  6. 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

      outCounts

      Número total de registros de dados consumidos pelo cliente SDK.

      outBytes

      Volume total de dados consumidos pelo cliente SDK, em bytes.

      outRps

      Quantidade de solicitações por segundo enviadas pelo cliente SDK para consumir dados.

      outBps

      Taxa de bits transmitidos por segundo durante o consumo de dados pelo cliente SDK.

      count

      Nenhum.

      inBytes

      Volume total de dados enviados pelo servidor DTS, em bytes.

      DStoreRecordQueue

      Tamanho da fila de cache de dados quando o servidor DTS envia dados.

      inCounts

      Número total de registros de dados enviados pelo servidor DTS.

      inRps

      Quantidade de solicitações enviadas pelo servidor DTS por segundo.

      inBps

      Taxa de bits transmitidos por segundo quando o servidor DTS envia dados.

      __dt

      Timestamp em que o cliente SDK recebe os dados, em milissegundos.

      DefaultUserRecordQueue

      Tamanho da fila de cache de dados após a serialização.