Todos os produtos
Search
Central de documentação

Data Transmission Service:Consumir dados assinados usando um SDK

Última atualização: Sep 14, 2026

Após configurar um canal de rastreamento de alterações com a criação de uma tarefa de rastreamento e um grupo de consumidores, use o kit de desenvolvimento de software (SDK) fornecido pelo DTS para consumir os dados assinados. Este tópico descreve como utilizar o código de exemplo.

Nota

Pré-requisitos

Precauções

  • Ao consumir dados assinados, conclua o consumo dos dados antes de chamar o método commit de DefaultUserRecord para confirmar as informações de offset. Não chame o método commit antes de consumir os dados, pois isso resultará em perda de dados.

  • Os processos de consumo são independentes entre si.

  • No console, Current Offset indica o offset que a tarefa de rastreamento assinou, e não o offset confirmado pelo cliente.

Procedimento

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

  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 no 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 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 os arquivos descompactados em um ambiente de desenvolvimento integrado (IDE).

    2. Abra o arquivo Java correspondente ao modo de uso desejado para o cliente SDK.

      Nota

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

      Modo de uso

      Arquivo Java

      Descrição

      Cenários

      Modo ASSIGN

      DTSConsumerAssignDemo.java

      Para garantir a ordem global das mensagens, o DTS atribui apenas uma partição (partição 0) a cada tópico de rastreamento. Se usar o cliente SDK no modo ASSIGN, inicie apenas um cliente SDK.

      Apenas um cliente SDK em um grupo de consumidores consome os dados assinados.

      Modo SUBSCRIBE

      DTSConsumerSubscribeDemo.java

      Para garantir a ordem global das mensagens, o DTS atribui apenas uma partição (partição 0) a cada tópico de rastreamento. No modo SUBSCRIBE, é possível iniciar vários clientes SDK no mesmo grupo de consumidores para recuperação de desastres. Se o cliente que está consumindo dados falhar, outro cliente SDK será atribuído automática e aleatoriamente à partição 0 para continuar o consumo.

      Vários clientes SDK no mesmo grupo de consumidores consomem dados assinados. Trata-se de um cenário de recuperação de desastres de dados.

    3. Defina os parâmetros no código Java.

      Código de exemplo

      ******        
          public static void main(String[] args) {
              // The Kafka broker URL.
              String brokerUrl = "dts-cn-***.com:18001";
              // The topic from which to consume data. The partition is 0.
              String topic = "cn_***_version2";
              // The username, password, and SID for authentication.
              String sid = "dts***";
              String userName = "dts***";
              String password = "DTS***";
              // The initial checkpoint for the first seek. This is a UNIX timestamp. For example, set this parameter to 1566180200 if you want to start consumption from 10:03:21 (CST) on Monday, August 19, 2019.
              String initCheckpoint = "1740472***";
              // If you use the SUBSCRIBE mode, you must configure the group. The Kafka consumer group is enabled.
              ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE;
        
              DTSConsumerSubscribeDemo consumerDemo = new DTSConsumerSubscribeDemo(brokerUrl, topic, sid, userName, password, initCheckpoint, subscribeMode);
              consumerDemo.start();
          }
      ******

      Parâmetro

      Descrição

      Como obter

      brokerUrl

      Especifica o endereço de rede e o número da porta do canal de rastreamento de alterações.

      Nota
      • Se o servidor onde você implanta o cliente SDK, como uma instância ECS, e a instância de rastreamento de alterações estiverem na mesma virtual private cloud (VPC), consuma dados pela VPC para reduzir a latência de rede.

      • 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 assinatura de destino. Na página Basic Information, obtenha o endereço de rede e o número da porta na seção Network.

      topic

      O tópico do canal de rastreamento de alterações.

      No console do DTS, clique em no ID da instância de assinatura de destino. Na página Basic Information, acesse a seção Basic Information e obtenha o Topic.

      sid

      O ID do grupo de consumidores.

      No console do DTS, clique em no ID da instância de assinatura de destino. Na página Consume Data, obtenha o Consumer Group ID/Name e a Account.

      userName

      O nome de usuário do grupo de consumidores.

      Aviso

      Se não utilizar o cliente fornecido neste tópico, defina o nome de usuário no formato <consumer group username>-<consumer group ID> (por exemplo, dtstest-dtsaebpv). Caso contrário, a conexão falhará.

      password

      A senha da conta.

      A senha definida para o nome de usuário do grupo de consumidores durante a criação do grupo.

      initCheckpoint

      O offset do consumidor. Representa o timestamp UNIX do primeiro registro de dados a ser consumido pelo cliente SDK, como 1620962769.

      Nota

      Utilize o offset do consumidor para:

      • Retomar o consumo de dados a partir de um offset específico após interrupção da aplicação, evitando perda de dados.

      • Ajustar o offset inicial para consumo de dados conforme necessário.

      O offset do consumidor deve estar dentro do intervalo de timestamps da instância de rastreamento e convertido para timestamp UNIX.

      Nota

      A coluna Data Range na lista de tarefas de rastreamento exibe o intervalo de timestamps da instância de assinatura de destino.

      subscribeMode

      O modo de utilização do cliente SDK. Não é necessário modifique este parâmetro.

      • ConsumerContext.ConsumerSubscribeMode.ASSIGN: Modo ASSIGN.

      • ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE: Modo SUBSCRIBE.

      N/A

  4. Abra a estrutura do projeto no seu IDE e certifique-se de que a versão do OpenJDK do projeto seja 1,8.

  5. Execute o código do cliente.

    Nota

    Na primeira execução do código, o IDE precisa de algum tempo para carregar automaticamente os plugins e dependências necessários.

    Resultados de exemplo (clique em para expandir)

    Resultados de execução normal

    Se o resultado abaixo for retornado, o cliente está funcionando corretamente e pronto para assinar alterações de dados do banco de dados de origem.

    ******
    [2025-02-25 18:47:22.991] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.consumer.KafkaConsumer:1587] - [Consumer clientId=consumer-dtsl5vy2ao5250****-1, groupId=dtsl5vy2ao5250****] Seeking to offset 8200 for partition cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0
    [2025-02-25 18:47:22.993] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ConsumerWrap:116] - RecordFetcher consumer:  subscribe for [cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0] with checkpoint [Checkpoint[ topicPartition: cn_hangzhou_vpc_rm_bp15uddebh4a1****_dts****_version2-0timestamp: 174048****, offset: 8200, info: 174048****]] start
    [2025-02-25 18:47:23.011] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174048044****,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:47:23.226] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8200]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [HEARTBEAT]
    
    [2025-02-25 18:47:23.226] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8201]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [HEARTBEAT]
    ******

    Resultados de assinatura normal

    Se o resultado abaixo for retornado, o cliente assinou com sucesso as alterações de dados (uma operação UPDATE) do banco de dados de origem.

    ******
    [2025-02-25 18:48:24.905] [INFO ] [com.aliyun.dts.subscribe.clients.recordprocessor.EtlRecordProcessor] [com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener:49] - 
    RecordID [8413]
    RecordTimestamp [174048****] 
    Source [{"sourceType": "MySQL", "version": "8.0.36"}]
    RecordType [UPDATE]
    Schema info [{, 
    recordFields= [{fieldName='id', rawDataTypeNum=8, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}, {fieldName='name', rawDataTypeNum=253, isPrimaryKey=false, isUniqueKey=false, fieldPosition=1}], 
    databaseName='dtsdb', 
    tableName='person', 
    primaryIndexInfo [[indexType=PrimaryKey, indexFields=[{fieldName='id', rawDataTypeNum=8, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}], cardinality=0, nullable=true, isFirstUniqueIndex=false, name=null]], 
    uniqueIndexInfo [[]], 
    partitionFields = null}]
    Before image {[Field [id] [3]
    Field [name] [test1]
    ]}
    After image {[Field [id] [3]
    Field [name] [test2]
    ]}
    ******

    Resultados de execução anormal

    Se o resultado abaixo for retornado, o cliente não consegue se conectar ao banco de dados de origem.

    ******
    [2025-02-25 18:22:18.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174047893****,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:22:22.002] [WARN ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.NetworkClient:780] - [Consumer clientId=consumer-dtsnd7u2n0625m****-1, groupId=dtsnd7u2n0625m****] Connection to node 1 (47.118.XXX.XXX/47.118.XXX.XXX:18001) could not be established. Broker may not be available.
    [2025-02-25 18:22:22.509] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ClusterSwitchListener:44] - Cluster not changed on update:5aPLLlDtTHqP8sKq-DZVfg
    [2025-02-25 18:22:23.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":1740478943160,"DefaultUserRecordQueue":0.0}
    [2025-02-25 18:22:27.192] [WARN ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [org.apache.kafka.clients.NetworkClient:780] - [Consumer clientId=consumer-dtsnd7u2n0625m****1, groupId=dtsnd7u2n0625m****] Connection to node 1 (47.118.XXX.XXX/47.118.XXX.XXX:18001) could not be established. Broker may not be available.
    [2025-02-25 18:22:27.618] [INFO ] [com.aliyun.dts.subscribe.clients.recordfetcher.KafkaRecordFetcher] [com.aliyun.dts.subscribe.clients.recordfetcher.ClusterSwitchListener:44] - Cluster not changed on update:5aPLLlDtTHqP8sKq-DZVfg
    ******

    O cliente SDK coleta e exibe periodicamente estatísticas sobre o consumo de dados. Essas estatísticas incluem o número total de registros de dados enviados e recebidos, o volume total de dados e o número de registros por segundo (RPS).

    [2025-02-25 18:22:18.160] [INFO ] [subscribe-logMetricsReporter-1-thread-1] [log.metrics:184] - {"outCounts":0.0,"outBytes":0.0,"outRps":0.0,"outBps":0.0,"count":11.0,"inBytes":0.0,"DStoreRecordQueue":0.0,"inCounts":0.0,"inRps":0.0,"inBps":0.0,"__dt":174047893****,"DefaultUserRecordQueue":0.0}

    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

    Número de solicitações por segundo para consumo de dados pelo cliente SDK.

    outBps

    Número de bits transmitidos por segundo durante o consumo de dados pelo cliente SDK.

    count

    Número total de parâmetros nas informações de consumo de dados (métricas).

    Nota

    Não inclui o próprio count.

    inBytes

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

    DStoreRecordQueue

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

    inCounts

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

    inBps

    Número de bits transmitidos por segundo quando o servidor DTS envia dados.

    inRps

    Número de solicitações 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 serialização.

  6. Edite o código para consumir os dados assinados conforme necessário.

    Ao consumir dados assinados, você precisa manage consumer offsets para evitar perda de dados, minimizar duplicação de dados e habilitar o consumo sob demanda.

Perguntas frequentes

  • O que fazer se eu não conseguir me conectar a uma instância de assinatura?

    Solucione o problema com base na mensagem de erro. Para mais informações, consulte Troubleshooting.

  • Qual é o formato de dados de um offset de consumidor após persistência?

    Após a persistência de um offset de consumidor, os dados são retornados no formato JSON. O offset persistido é um timestamp Unix que pode ser passado diretamente ao SDK. Por exemplo, nos dados retornados, o valor 1700709977 para a chave "timestamp" representa o offset de consumidor persistido.

    {"groupID":"dtsglg11d48230***","streamCheckpoint":[{"partition":0,"offset":577989,"topic":"ap_southeast_1_vpc_rm_t4n22s21iysr6****_root_version2","timestamp":170070****,"info":""}]}
  • Uma tarefa de rastreamento pode ser consumida por vários clientes em paralelo?

    Não. Embora o modo SUBSCRIBE permita que vários clientes sejam executados em paralelo, apenas um cliente pode consumir dados por vez.

  • Qual versão do cliente Kafka está encapsulada no código do SDK?

    As versões 2.0.0 e posteriores do dts-new-subscribe-sdk encapsulam o Kafka client (kafka-clients) 2.7.0. Versões anteriores a 2.0.0 encapsulam o Kafka client 1.0.0.

    Nota

    Se utilizar uma ferramenta de detecção de vulnerabilidades em pacotes de dependência dependency package vulnerability detection tool no desenvolvimento da sua aplicação e descobrir que o Kafka client (kafka-clients) encapsulado pelo dts-new-subscribe-sdk possui uma vulnerabilidade de segurança, resolva essa vulnerabilidade substituindo o cliente pela versão 2.1.4-shaded.

    <dependency>
        <groupId>com.aliyun.dts</groupId>
        <artifactId>dts-new-subscribe-sdk</artifactId>
        <version>2.1.4-shaded</version>
    </dependency>

Apêndice

Gerenciar offsets de consumidor

Quando o cliente SDK inicia pela primeira vez, reinicia ou tenta nova conexão internamente, consulte e passe um offset de consumidor para iniciar ou retomar o consumo de dados. O offset de consumidor corresponde ao timestamp UNIX do primeiro registro de dados que o cliente SDK consumirá.

Para redefinir o offset de consumidor do cliente, consulte e modifique o offset com base no modo de consumo (modo de uso do SDK), conforme descrito na tabela a seguir.

Nota

Os níveis de armazenamento de checkpoint listados acima (armazenamento externo, arquivo localCheckpointStore e offset do DTS Server/DStore) são gravados em uma única operação de confirmação, portanto seus conteúdos são quase idênticos. A principal diferença é que o arquivo localCheckpointStore será perdido se o nó ou contêiner for destruído, enquanto o offset do DTS Server (DStore) persiste no lado do servidor.

Cenário

Modo de uso do SDK

Método de gerenciamento de offset

Consultar um offset de consumidor

Modo ASSIGN, Modo SUBSCRIBE

  • Como o cliente SDK salve o offset da mensagem a cada 5 segundos e o confirma no servidor DTS, consulte o último offset de consumidor em um dos seguintes locais:

    • Arquivo localCheckpointStore no servidor onde o cliente SDK está localizado.

    • Página Data Consumption do canal de rastreamento de alterações.

  • Caso tenha configurado um meio de armazenamento compartilhado persistente externo, como um banco de dados, no arquivo consumerContext.java usando setUserRegisteredStore(new UserMetaStore()), esse meio de armazenamento salve o offset da mensagem a cada 5 segundos para consulta.

O cliente SDK inicia pela primeira vez. Passe um offset de consumidor para consumir dados.

Modo ASSIGN, Modo SUBSCRIBE

Dependendo do padrão de uso do cliente SDK, selecione o arquivo Java DTSConsumerAssignDemo.java ou DTSConsumerSubscribeDemo.java e configure the consumer offset (initCheckpoint) para consumir dados.

O cliente SDK tenta nova conexão internamente. Passe o último offset de consumidor registrado para retomar o consumo de dados.

Modo ASSIGN

Pesquise o último offset de consumidor registrado na seguinte ordem. Se um offset for encontrado, as informações de offset serão retornadas:

  1. O meio de armazenamento externo configurado usando setUserRegisteredStore(new UserMetaStore()) no arquivo consumerContext.java.

  2. O arquivo localCheckpointStore no servidor onde o cliente SDK está localizado.

  3. O offset salvo no DTS Server (DStore).

  4. O timestamp inicial passado para initCheckpoint no arquivo DTSConsumerSubscribeDemo.java.

Modo SUBSCRIBE

Pesquise o último offset de consumidor registrado na seguinte ordem. Se um offset for encontrado, as informações de offset serão retornadas:

  1. O meio de armazenamento externo configurado no arquivo consumerContext.java usando setUserRegisteredStore(new UserMetaStore()).

  2. O offset salvo no DTS Server (módulo de ingestão de dados incrementais).

    Nota

    Este offset é atualizado somente após o cliente SDK chamar o método commit para atualize o offset de consumidor.

  3. O timestamp inicial passado para initCheckpoint no arquivo DTSConsumerSubscribeDemo.java.

  4. O offset inicial do DTS Server (novo módulo de ingestão de dados incrementais).

    Importante

    Se o módulo de ingestão de dados incrementais for alternado, o novo módulo não poderá salve o último offset de consumidor do cliente. Isso pode fazer com que o consumo de dados comece a partir de um offset mais antigo. Recomendamos que você persistently store the consumer offset no cliente.

O cliente SDK é reiniciado. Passe o último offset de consumidor registrado para retomar o consumo de dados.

Modo ASSIGN

Com base na configuração setForceUseCheckpoint no arquivo consumerContext.java, o offset de consumidor é consultado e, se encontrado, as informações de offset são retornadas:

  • Quando definido como true, o cliente SDK usará o initCheckpoint passado como offset de consumidor sempre que for reiniciado.

  • Quando configurado como false ou não configurado, localize o offset de consumidor do registro anterior na seguinte ordem:

    1. O meio de armazenamento externo configurado no arquivo consumerContext.java usando setUserRegisteredStore(new UserMetaStore()).

    2. O arquivo localCheckpointStore no servidor onde o cliente SDK está localizado.

    3. O offset salvo no DTS Server (módulo de ingestão de dados incrementais).

      Nota

      Este offset é atualizado somente após o cliente SDK chamar o método commit para atualize o offset de consumidor.

Modo SUBSCRIBE

Neste modo, no arquivo consumerContext.java, a configuração setForceUseCheckpoint não tem efeito. Localize o offset de consumidor do registro anterior na seguinte ordem:

  1. O meio de armazenamento externo configurado no arquivo consumerContext.java com setUserRegisteredStore(new UserMetaStore()).

  2. O offset salvo no DTS Server (módulo de ingestão de dados incrementais).

    Nota

    Este offset é atualizado somente após o cliente SDK chamar o método commit para atualize o offset de consumidor.

  3. O timestamp inicial passado para initCheckpoint no arquivo DTSConsumerSubscribeDemo.java.

  4. O offset inicial do DTS Server (novo módulo de ingestão de dados incrementais).

Armazenar o offset de consumidor com persistência

Se ocorrer um failover de recuperação de desastres para o módulo de ingestão de dados incrementais, o novo módulo não poderá salve o último offset de consumidor do cliente. Isso é especialmente verdadeiro no modo SUBSCRIBE. Como resultado, o cliente pode começar a consumir dados a partir de um offset anterior, levando ao consumo repetido de dados históricos. Por exemplo, suponha que, antes de um failover de serviço, o intervalo de offset do módulo antigo seja de 08:00:00 em 11 de novembro de 2023 a 08:00:00 em 12 de novembro de 2023, e o offset de consumidor do cliente seja 08:00:00 em 12 de novembro de 2023. Após o failover, o intervalo de offset do novo módulo vai de 10:00:00 em 08 de novembro de 2023 a 08:01:00 em 12 de novembro de 2023. Nesse cenário, o cliente inicia o consumo a partir do offset inicial do novo módulo (10:00:00 em 08 de novembro de 2023), o que resulta em consumo repetido de dados.

Para evitar o consumo repetido de dados históricos nesse cenário de failover, recomendamos configure um método de armazenamento persistente para o offset de consumidor no cliente. O método de exemplo abaixo é fornecido como referência e pode ser modificado conforme necessário.

  1. Crie um método UserMetaStore() que herde e implemente o método AbstractUserMetaStore().

    Por exemplo, utilize um banco de dados MySQL para armazenar informações de offset. O código Java de exemplo abaixo demonstra como fazer isso:

    public class UserMetaStore extends AbstractUserMetaStore {
    
        @Override
        protected void saveData(String groupID, String toStoreJson) {
            Connection con = getConnection();
            String sql = "insert into dts_checkpoint(group_id, checkpoint) values(?, ?)";
    
            PreparedStatement pres = null;
            ResultSet rs = null;
    
            try {
                pres = con.prepareStatement(sql);
                pres.setString(1, groupID);
                pres.setString(2, toStoreJson);
                pres.execute();
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                close(rs, pres, con);
            }
        }
    
        @Override
        protected String getData(String groupID) {
            Connection con = getConnection();
            String sql = "select checkpoint from dts_checkpoint where group_id = ?";
    
            PreparedStatement pres = null;
            ResultSet rs = null;
    
            try {
                pres = con.prepareStatement(sql);
                pres.setString(1, groupID);
                rs = pres.executeQuery();
                              
                if (rs.next()) {
                    String checkpoint = rs.getString("checkpoint");
                    return checkpoint;
                }
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                close(rs, pres, con);
            }
            return null;
        }
    }
    
  2. No arquivo consumerContext.java, chame o método setUserRegisteredStore(new UserMetaStore()) para configure o meio de armazenamento externo.

Solução de problemas

Exceção

Mensagem de erro

Causa

Solução

Falha na conexão

ERROR
CheckResult{isOk=false, errMsg='telnet dts-cn-hangzhou.aliyuncs.com:18009
failed, please check the network and if the brokerUrl is correct'}
(com.aliyun.dts.subscribe.clients.DefaultDTSConsumer)

O brokerUrl está incorreto.

Insira o brokerUrl, userName e password corretos. Para mais informações, consulte Parameter descriptions.

telnet real node *** failed, please check the network

Não é possível conectar ao endereço IP real usando o endereço do broker.

ERROR CheckResult{isOk=false, errMsg='build kafka consumer failed, error: org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata, probably the user name or password is wrong'} (com.aliyun.dts.subscribe.clients.DefaultDTSConsumer)

O nome de usuário ou a senha estão incorretos.

com.aliyun.dts.subscribe.clients.exception.TimestampSeekException: RecordGenerator:seek timestamp for topic [cn_hangzhou_rm_bp11tv2923n87081s_rdsdt_dtsacct-0] with timestamp [1610249501] failed

No arquivo consumerContext.java, setUseCheckpoint está definido como true, mas o offset de consumidor não está dentro do intervalo de timestamps da instância de assinatura.

Insira um offset de consumidor dentro do intervalo de timestamps da instância de assinatura. Para mais informações sobre o método de consulta, consulte Parameter description.

Lentidão no consumo da assinatura

N/A

  • Analise o motivo da lentidão no consumo de dados consultando os parâmetros em statistics information referentes ao tamanho das filas DStoreRecordQueue e DefaultUserRecordQueue.

    • Se o valor de DStoreRecordQueue permanecer 0, a velocidade de extração de dados pelo servidor DTS está lenta.

    • Se o valor de DefaultUserRecordQueue permanecer no valor padrão de 512, a velocidade de consumo de dados pelo cliente SDK está lenta.

  • Redefina o offset modificando o offset de consumidor (initCheckpoint) no código conforme necessário.