Todos os produtos
Search
Central de documentação

Data Transmission Service:Usar uma demonstração de SDK para consumir dados de alteração do PolarDB-X 1.0

Última atualização: Jun 27, 2026

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.

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

  2. Crie um ou mais grupos de consumidores. Para mais informações, consulte Criar um grupo de consumidores.

  3. Baixe e descompacte a demonstração de SDK. Para obter o link de download, consulte Código da Demonstração de SDK.

    Importante

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

  4. Abra o projeto no IntelliJ IDEA.

    1. Abra o IntelliJ IDEA e clique em Open or Import.

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

    3. Na caixa de diálogo exibida, selecione Open as Project.

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

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

    Nota

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

    Nota

    Use um mecanismo de busca para encontrar um conversor de timestamp Unix.

  7. Na barra de menu superior do IntelliJ IDEA, escolha Run > Run para execute o cliente.

    Nota

    Na 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

      outCounts

      Quantidade total de registros de dados consumidos pelo cliente SDK.

      outBytes

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

      outRps

      Taxa de consumo de dados do cliente SDK, em requisições por segundo (RPS).

      outBps

      Taxa de consumo de dados do cliente SDK, em bits por segundo (bps).

      count

      Parâmetro reservado.

      inBytes

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

      DStoreRecordQueue

      Tamanho atual da fila de cache de dados no servidor DTS.

      inCounts

      Quantidade total de registros de dados enviados pelo servidor DTS.

      inRps

      Taxa de envio de dados do servidor DTS, em requisições por segundo (RPS).

      inBps

      Taxa de envio de dados do servidor DTS, em bits por segundo (bps).

      __dt

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

      DefaultUserRecordQueue

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

  8. 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);
        }