Todos os produtos
Search
Central de documentação

Data Transmission Service:Use the SDK demo to consume the data tracked from a PolarDB-X 1.0 instance

Última atualização: Jun 27, 2026

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:

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

  1. Abra o IntelliJ IDEA. Na janela de boas-vindas, clique em Open or Import.

    Open a project

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

    1

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

Find the Java file

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

accessKeyId

accesskey ID

Consulte Criar e obter um par de AccessKeys.

accessKeySecret

AccessKey secret

Consulte Criar e obter um par de AccessKeys.

regionId

ID da região da tarefa de rastreamento de alterações. Exemplo: cn-hangzhou para China (Hangzhou).

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.

dtsInstanceId

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.

jobId

ID da tarefa de rastreamento de alterações.

Chame a operação da API DescribeDtsJobs para consultar o ID da tarefa.

sid

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.

userName

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.

password

Senha da conta do grupo de consumidores.

A senha é defina durante a criação do grupo de consumidores.

proxyUrl

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.

checkpoint

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.

Nota

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

outCounts

Número total de registros de dados consumidos pelo cliente do sdk

outBytes

Quantidade total de dados consumidos pelo cliente do sdk, em bytes

outRps

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

outBps

Throughput de consumo do cliente do sdk, em bits por segundo

count

Nenhum

inBytes

Quantidade total de dados enviados pelo servidor DTS, em bytes

DStoreRecordQueue

Tamanho da fila de cache de dados no lado do servidor DTS

inCounts

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

inRps

Taxa de envio do servidor DTS, em RPS

inBps

Throughput de envio do servidor DTS, em bits por segundo

__dt

Timestamp em que o cliente do sdk recebeu os dados, em milissegundos

DefaultUserRecordQueue

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

Próximos passos