Após configurar uma instância de Data Subscription, use o código de exemplo do SDK fornecido pelo Data Transmission Service (DTS) para consumir os dados de alteração.
Procedimento
Se a fonte de dados for uma instância PolarDB-X 1.0 ou um banco de dados lógico do DMS, consulte Consumir dados de Data Subscription do PolarDB-X 1.0 usando código de exemplo do SDK.
Se você usar um usuário RAM para consumir dados, esse usuário deverá ter a permissão AliyunDTSFullAccess e as permissões necessárias para acessar os objetos da assinatura. Para obter mais informações sobre como conceder permissões, consulte Autorizar um usuário RAM a gerenciar instâncias do DTS usando uma política do sistema e Gerenciar permissões de usuário RAM.
Cada consumidor opera de forma independente.
Este tópico fornece um cliente SDK de exemplo em Java. Para códigos de exemplo em Python e Go, consulte dts-subscribe-demo.
O procedimento a seguir demonstra como executar o código de exemplo do SDK para consumir dados de Data Subscription no IntelliJ IDEA (Community Edition 2020.1 para Windows).
Crie uma instância de Data Subscription. Para mais detalhes, consulte Criar um canal de Data Subscription para uma instância ApsaraDB RDS for MySQL, Criar um canal de Data Subscription para um cluster PolarDB for MySQL ou Criar um canal de Data Subscription para um banco de dados Oracle.
-
Crie um ou mais grupos de consumidores. Para mais informações, consulte Criar grupos de consumidores.
ImportanteAo consumir dados de Data Subscription, chame o método
commitdeDefaultUserRecordpara confirmar os checkpoints. A ausência dessa confirmação pode resultar em consumo duplicado de dados. -
Use o código de exemplo do SDK conforme as necessidades do seu negócio.
-
Usar o novo pacote SDK de Data Subscription (recomendado)
Abra o IntelliJ IDEA e clique em Create New Project para criar um projeto para sua aplicação.
No projeto, localize o arquivo de modelo de objeto do projeto (POM): pom.xml.
-
Adicione a seguinte dependência ao arquivo pom.xml:
<dependency> <groupId>com.aliyun.dts</groupId> <artifactId>dts-new-subscribe-sdk</artifactId> <version>{dts_new_sdk_version}</version> </dependency>NotaEncontre a dependência Maven mais recente na página dts-new-subscribe-sdk.
Para mais informações sobre como usar o novo SDK de assinatura, consulte Usar o código de exemplo.
-
Usar uma versão personalizada do novo SDK de Data Subscription
-
Baixe o pacote de código de exemplo do SDK e descompacte-o.
NotaClique em
e selecione Download ZIP para baixar o pacote. -
Acesse o diretório do código de exemplo do SDK descompactado. Abra o arquivo pom.xml com um editor de texto e atualize o SDK de Data Subscription para a versão mais recente.
ImportanteObtenha a versão mais recente do SDK de Data Subscription no site do Maven. Para mais detalhes, consulte a página do Maven para o SDK de Data Subscription.
Abra o IntelliJ IDEA e clique em Open or Import.

Na caixa de diálogo exibida, acesse o diretório do código de exemplo do SDK descompactado, expanda as pastas e localize o arquivo pom.xml.

Na caixa de diálogo exibida, selecione Open as Project.
-
No IntelliJ IDEA, expanda as pastas. De acordo com o modo de uso do cliente SDK, selecione e clique duas vezes no arquivo Java correspondente: DTSConsumerAssignDemo.java ou DTSConsumerSubscribeDemo.java.
NotaO DTS oferece suporte aos seguintes modos de uso do cliente SDK:
Modo ASSIGN: Para garantir a ordem global das mensagens, o DTS atribui apenas uma partição (partição 0) a cada tópico de assinatura. Ao usar um cliente SDK no modo ASSIGN, inicie apenas um cliente.
Modo SUBSCRIBE: Para assegurar a ordem global das mensagens, o DTS designa somente uma partição (partição 0) por tópico de assinatura. Se você usar o cliente SDK no modo SUBSCRIBE, poderá iniciar vários clientes SDK em um grupo de consumidores para recuperação de desastres. Caso o cliente ativo falhe, outro cliente SDK será automaticamente atribuído à partição 0 para retomar o consumo.
-
-
-
Defina os parâmetros necessários no arquivo Java.

Tabela 1. Parâmetros obrigatórios
Parâmetro
Descrição
Origem
brokerUrlO endpoint e o número da porta da instância de Data Subscription.
Nota-
Se a instância ECS que executa o cliente SDK e a instância de Data Subscription estiverem na mesma rede clássica ou Virtual Private Cloud (VPC), use o endpoint interno para a assinatura visando minimizar 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 ID da instância de Data Subscription desejada. Na página Basic Information, obtenha o endpoint e o número da porta na seção Network.
topicO tópico de assinatura da instância.
No console do DTS, clique em ID da instância de Data Subscription alvo. Na página Basic Information, obtenha o Topic na seção Basic Information.
sidO ID do grupo de consumidores.
No console do DTS, clique em ID da instância de Data Subscription correspondente e, em seguida, clique em Consume Data. Obtenha o Consumer Group ID e a Account do grupo de consumidores.
NotaA senha do nome de usuário do grupo de consumidores é definida durante a criação do grupo.
userNameO nome de usuário do grupo de consumidores.
AvisoCaso não use o cliente fornecido neste tópico, defina o nome de usuário no formato
<Username>-<Consumer Group ID>. Exemplo:dtstest-dtsaebpv. Caso contrário, a conexão falhará.passwordA senha do nome de usuário.
initCheckpointO checkpoint de consumo, especificado como um timestamp UNIX, a partir do qual o cliente SDK começa a consumir dados. Exemplo: 1620962769.
NotaAs informações de checkpoint de consumo são úteis nos seguintes cenários:
-
Para retomar o consumo e evitar perda de dados após uma interrupção da aplicação, informe o último checkpoint de consumo conhecido.
-
Ao iniciar o cliente, passe um checkpoint de consumo específico para consumir dados a partir de uma posição desejada.
O checkpoint de consumo deve estar dentro do intervalo de dados da instância de Data Subscription (conforme ilustrado na figura) e convertido para um timestamp UNIX.
NotaUse um mecanismo de busca para encontrar um conversor de timestamp UNIX.
ConsumerContext.ConsumerSubscribeMode subscribeModeO modo de uso do cliente SDK. Valores válidos:
-
ConsumerContext.ConsumerSubscribeMode.ASSIGN: Modo ASSIGN. Apenas um cliente SDK em um grupo de consumidores pode consumir dados de Data Subscription. -
ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE: Modo SUBSCRIBE. Permite iniciar múltiplos clientes SDK no mesmo grupo de consumidores para fins de recuperação de desastres.
N/A
-
-
Na barra de navegação superior do IntelliJ IDEA, escolha para executar o cliente.
NotaNa primeira execução do cliente, o carregamento e a instalação automática das dependências necessárias podem levar algum tempo.
A figura a seguir exibe o resultado. O cliente consome com êxito os dados de alteração do banco de dados de origem.

-
Periodicamente, o cliente SDK agrega e exibe estatísticas sobre o consumo de dados, incluindo a contagem total e o volume de registros enviados e recebidos, além de solicitações por segundo (RPS).

Tabela 2. Estatísticas de consumo de dados
Parâmetro
Descrição
outCountsQuantidade total de registros de dados consumidos pelo cliente SDK.
outBytesVolume total de dados consumidos pelo cliente SDK. Unidade: bytes.
outRpsNúmero de solicitações por segundo (RPS) com que o cliente SDK consome dados.
outBpsTaxa de consumo de dados do cliente SDK, em bits por segundo (bps).
inBytesVolume total de dados enviados pelo servidor DTS. Unidade: bytes.
DStoreRecordQueueTamanho da fila interna de cache de dados para registros recebidos do servidor DTS.
inCountsQuantidade total de registros de dados enviados pelo servidor DTS.
inRpsRPS com que o servidor DTS envia dados.
__dtTimestamp em que o cliente SDK recebe os dados. Unidade: milissegundos.
DefaultUserRecordQueueTamanho da fila de dados que contém registros prontos para processamento pela aplicação consumidora.
Salvar e consultar checkpoints de consumo
Para iniciar ou retomar o consumo de dados (por exemplo, na primeira execução, reinicialização ou nova tentativa interna), o cliente SDK requer um checkpoint de consumo. A tabela abaixo descreve como gerenciar e consultar checkpoints em diferentes cenários para evitar perda de dados, minimizar o consumo duplicado e permitir o consumo sob demanda.
|
Cenário |
Modo de uso do SDK |
Método de consulta |
|
Consultar um checkpoint de consumo |
Modo ASSIGN, Modo SUBSCRIBE |
|
|
Primeira inicialização: passar um checkpoint para iniciar o consumo. |
Modo ASSIGN, Modo SUBSCRIBE |
De acordo com o modo de uso do cliente SDK, selecione o arquivo DTSConsumerAssignDemo.java ou DTSConsumerSubscribeDemo.java e configure o parâmetro |
|
O cliente SDK precisa passar novamente o último checkpoint de consumo registrado para continuar o consumo após uma nova tentativa interna. |
Modo ASSIGN |
Busque o último checkpoint de consumo registrado na seguinte ordem. A pesquisa para e retorna as informações do checkpoint assim que ele é encontrado:
|
|
Modo SUBSCRIBE |
Pesquise o último checkpoint de consumo registrado seguindo a ordem abaixo. A busca é interrompida e retorna as informações assim que o checkpoint é localizado:
|
|
|
O cliente SDK foi reiniciado e precisa passar novamente o último checkpoint de consumo registrado para prosseguir com o consumo. |
Modo ASSIGN |
Consulte o checkpoint de consumo com base na configuração
|
|
Modo SUBSCRIBE |
Neste modo, a configuração
|
Persistir checkpoints de consumo
Durante um evento de recuperação de desastres no módulo de coleta de dados incrementais (especialmente no modo SUBSCRIBE), o novo módulo não retém o checkpoint de consumo mais recente do cliente. O cliente pode retomar a partir de um checkpoint mais antigo, resultando em consumo duplicado de dados históricos. Por exemplo, antes da troca, o intervalo de checkpoints do módulo antigo vai de 08:00:00 em 11 de novembro de 2023 até 08:00:00 em 12 de novembro de 2023, e o checkpoint do cliente é 08:00:00 em 12 de novembro de 2023. Após a troca, o intervalo de checkpoints do novo módulo vai de 10:00:00 em 08 de novembro de 2023 até 08:01:00 em 12 de novembro de 2023. O cliente inicia a partir do checkpoint inicial do novo módulo (10:00:00 em 08 de novembro de 2023), causando consumo duplicado.
Para evitar consumo duplicado nesse cenário, configure um armazenamento persistente de checkpoints no cliente. O exemplo a seguir apresenta uma implementação possível que você pode adaptar às suas necessidades.
-
Crie uma classe
UserMetaStoreque estendaAbstractUserMetaStore.Por exemplo, para armazenar informações de checkpoint em um banco de dados MySQL, use o seguinte código Java:
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); ResultSet rs = pres.executeQuery() String checkpoint = rs.getString("checkpoint"); return checkpoint; } catch (Exception e) { e.printStackTrace(); } finally { close(rs, pres, con); } } } No arquivo consumerContext.java, configure o meio de armazenamento externo utilizando o método
setUserRegisteredStore(new UserMetaStore()).
Perguntas frequentes
-
Como resolver problemas de conexão com uma instância de Data Subscription?
Solucione o problema com base na mensagem de erro. Para mais detalhes, consulte Solução de problemas.
-
Em qual formato os checkpoints de consumo são persistidos?
Os dados de checkpoint de consumo persistidos são armazenados no formato JSON. O checkpoint persistido é um timestamp UNIX que pode ser passado diretamente ao SDK. No exemplo de resposta abaixo, o valor
1700709977para a chave"timestamp"representa o checkpoint de consumo persistido.{"groupID":"dtsglg11d48230***","streamCheckpoint":[{"partition":0,"offset":577989,"topic":"ap_southeast_1_vpc_rm_t4n22s21iysr6****_root_version2","timestamp":1700709977,"info":""}]}
Solução de problemas
|
Problema |
Mensagem de erro |
Causa |
Solução |
|
Impossível conectar |
|
O |
Insira os valores corretos para os parâmetros |
|
O endereço do broker não consegue se conectar ao endereço IP real. |
||
|
O nome de usuário ou a senha está incorreta. |
||
|
No arquivo consumerContext.java, o parâmetro |
Informe um checkpoint de consumo que esteja dentro do intervalo de dados da instância de Data Subscription. Para mais detalhes, consulte Parâmetros obrigatórios. |
|
|
O consumo fica lento |
N/A |
|
|