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.
Para consumir dados assinados de uma fonte de dados PolarDB-X 1.0, consulte Use an SDK to consume subscribed data from PolarDB-X 1.0.
Este tópico fornece código de exemplo para um cliente SDK em Java. Para exemplos em Python e Go, consulte dts-subscribe-demo.
Pré-requisitos
-
Uma instância de assinatura criada e em execução no estado Normal.
NotaPara obter instruções sobre como criar uma instância de assinatura, consulte Subscription Plan Overview.
Você já created a consumer group para sua instância de assinatura.
Caso utilize um usuário RAM para consumir dados assinados, esse usuário deve ter a permissão AliyunDTSFullAccess e permissões de acesso aos objetos assinados. Para mais informações, consulte Grant permissions to a RAM user to manage DTS using a system policy e Manage the permissions of a RAM user.
Precauções
Ao consumir dados assinados, conclua o consumo dos dados antes de chamar o método
commitdeDefaultUserRecordpara 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
Baixe os arquivos de código do SDK de exemplo e descompacte o pacote.
-
Verifique a versão do código do SDK.
Acesse o diretório onde você descompactou o código de exemplo do SDK.
Abra o arquivo pom.xml no diretório com um editor de texto.
-
Atualize o SDK de rastreamento de alterações para a versão mais recente.
NotaEncontre a dependência Maven mais recente na página dts-new-subscribe-sdk.
-
Edite o código do SDK.
Abra os arquivos descompactados em um ambiente de desenvolvimento integrado (IDE).
-
Abra o arquivo Java correspondente ao modo de uso desejado para o cliente SDK.
NotaO 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.
-
Defina os parâmetros no código Java.
Parâmetro
Descrição
Como obter
brokerUrlEspecifica o endereço de rede e o número da porta do canal de rastreamento de alterações.
NotaSe 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.
topicO 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.
sidO 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.
userNameO nome de usuário do grupo de consumidores.
AvisoSe 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á.passwordA senha da conta.
A senha definida para o nome de usuário do grupo de consumidores durante a criação do grupo.
initCheckpointO offset do consumidor. Representa o timestamp UNIX do primeiro registro de dados a ser consumido pelo cliente SDK, como 1620962769.
NotaUtilize 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.
NotaA coluna Data Range na lista de tarefas de rastreamento exibe o intervalo de timestamps da instância de assinatura de destino.
subscribeModeO 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
Abra a estrutura do projeto no seu IDE e certifique-se de que a versão do OpenJDK do projeto seja 1,8.
-
Execute o código do cliente.
NotaNa primeira execução do código, o IDE precisa de algum tempo para carregar automaticamente os plugins e dependências necessários.
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
outCountsNúmero total de registros de dados consumidos pelo cliente SDK.
outBytesVolume total de dados consumidos pelo cliente SDK, em bytes.
outRpsNúmero de solicitações por segundo para consumo de dados pelo cliente SDK.
outBpsNúmero de bits transmitidos por segundo durante o consumo de dados pelo cliente SDK.
countNúmero total de parâmetros nas informações de consumo de dados (métricas).
NotaNão inclui o próprio
count.inBytesVolume total de dados enviados pelo servidor DTS, em bytes.
DStoreRecordQueueTamanho atual da fila de cache de dados quando o servidor DTS envia dados.
inCountsNúmero total de registros de dados enviados pelo servidor DTS.
inBpsNúmero de bits transmitidos por segundo quando o servidor DTS envia dados.
inRpsNúmero de solicitações por segundo quando o servidor DTS envia dados.
__dtTimestamp em que o cliente SDK recebe os dados, em milissegundos.
DefaultUserRecordQueueTamanho da fila de cache de dados após serialização.
-
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
1700709977para 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.
NotaSe 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.
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 |
|
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 ( |
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:
|
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:
| |
O cliente SDK é reiniciado. Passe o último offset de consumidor registrado para retomar o consumo de dados. | Modo ASSIGN | Com base na configuração
|
Modo SUBSCRIBE | Neste modo, no arquivo consumerContext.java, a configuração
|
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.
-
Crie um método
UserMetaStore()que herde e implemente o métodoAbstractUserMetaStore().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; } } 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 | | O | Insira o |
| Não é possível conectar ao endereço IP real usando o endereço do broker. | ||
| O nome de usuário ou a senha estão incorretos. | ||
| No arquivo consumerContext.java, | 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 |
| |