Este tópico descreve como usar a demonstração de um cliente Kafka para consumir dados rastreados. O recurso de rastreamento de alterações da nova versão permite consumir dados rastreados com um cliente Kafka das versões V0.11 a V2.7.
Observações de uso
-
Se você ativar o commit automático no recurso de rastreamento de alterações, alguns dados poderão ser confirmados antes do consumo, causando perda de informações. Recomendamos confirmar os dados manualmente.
NotaSe a confirmação dos dados falhar, reinicie o cliente para retomar o consumo a partir do último checkpoint de consumo registrado. No entanto, dados duplicados podem ser gerados nesse período. Filtre manualmente esses dados duplicados.
-
Os dados são serializados e armazenados no formato Avro. Para mais detalhes, consulte Record.avsc.
AvisoCaso não utilize o cliente Kafka descrito neste tópico, analise os dados rastreados com base no esquema Avro (Exemplo de desserialização Avro do DTS) e valide os dados analisados.
A unidade de busca é segundos quando o Data Transmission Service (DTS) chama a operação
offsetForTimes. Para um cliente Kafka nativo, a unidade é milissegundos ao chamar essa mesma operação.Conexões transitórias podem ocorrer entre o cliente Kafka e o servidor de rastreamento de alterações devido a diversos fatores, como recuperação de desastres. Se você não estiver usando o cliente Kafka descrito neste tópico, seu cliente deve ter capacidade de reconexão de rede.
Ao usar um cliente Kafka nativo para consumir dados rastreados, o módulo de coleta de dados incrementais pode sofrer alterações no DTS. No modo subscribe, o checkpoint de consumo salvo pelo cliente Kafka no servidor DTS é removido. Especifique um checkpoint de consumo para consumir os dados rastreados conforme suas necessidades de negócio. Para consumir dados no modo subscribe, use a demonstração do SDK fornecida pelo DTS para rastrear e consumir dados ou gerencie manualmente o checkpoint de consumo. Para mais informações, consulte Consume subscribed data using an SDK e a seção Manage the consumption checkpoint deste tópico.
Executar o cliente Kafka
Baixe a demonstração do cliente Kafka. Para mais informações sobre como usar a demonstração, consulte Readme.
Clique em
e selecione Download ZIP para baixar o pacote.Se usar um cliente Kafka versão 2,0, altere o número da versão no arquivo subscribe_example-master/javaimpl/pom.xml para 2.0.0.
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.0.0</version>
</dependency>
|
Etapa |
Diretório ou arquivo relacionado |
|
1. Use o consumidor nativo do Kafka para obter dados incrementais da instância de rastreamento de alterações. |
subscribe_example-master/javaimpl/src/main/java/recordgenerator/ |
|
2. Desserialize a imagem dos dados incrementais e obtenha a pré-imagem , a pós-imagem e outros atributos. Aviso
|
subscribe_example-master/javaimpl/src/main/java/boot/RecordPrinter.java |
|
3. Converta os valores dataTypeNumber nos dados desserializados para os tipos de dados do banco de dados correspondente. Nota
Para mais informações, consulte as seguintes seções deste tópico: |
subscribe_example-master/javaimpl/src/main/java/recordprocessor/mysql/ |
Procedimento
As etapas a seguir mostram como executar o cliente Kafka para consumir dados rastreados. Neste exemplo, usa-se o IntelliJ IDEA Community Edition 2018.1.4 para Windows.
Crie uma instância de rastreamento de alterações. Para mais informações, consulte Change tracking overview.
Crie um ou mais grupos de consumidores. Para mais informações, consulte Add a consumer group.
-
Baixe o pacote da demonstração do cliente Kafka e descompacte-o.
NotaClique em
e selecione Download ZIP para baixar o pacote. Abra o IntelliJ IDEA. Na janela exibida, clique em Open.
-
Na caixa de diálogo exibida, navegue até o diretório onde a demonstração baixada está localizada. Localize o arquivo pom.xml.
Navegue até
kafkademo>subscribe_example-master>javaimpl, selecionepom.xmle clique em OK. Na caixa de diálogo exibida, selecione Open as Project.
Na janela de ferramentas Project do IntelliJ IDEA, navegue pelas pastas para localizar o arquivo de demonstração do cliente Kafka e dê um duplo clique nele. O nome do arquivo é NotifyDemoDB.java.
-
Especifique os parâmetros no arquivo NotifyDemoDB.java.
public static Properties getConfigs() { Properties properties = new Properties(); // user password and sid for auth properties.setProperty(USER_NAME, "dtstest"); properties.setProperty(PASSWORD_NAME, "xxx"); properties.setProperty(SID_NAME, "dtsxxx"); // kafka consumer group general same with sid properties.setProperty(GROUP_NAME, "dtsxxx"); // topic to consume, partition is 0 properties.setProperty(KAFKA_TOPIC, "cn_hangzhou_xxx"); // kafka broker url properties.setProperty(KAFKA_BROKER_URL_NAME, "dts-cn-xxx.com:18001"); // initial checkpoint for first seek(a timestamp to set, eg 1566180200 if you want (Mon Aug 19 10:03:21 CST 2019)) properties.setProperty(INITIAL_CHECKPOINT_NAME, "1583307907"); // if force use config checkpoint when start. for checkpoint reset properties.setProperty(USE_CONFIG_CHECKPOINT_NAME, "true"); // use consumer assign or subscribe interface // when use subscribe mode, group config is required. kafka consumer group is enabled properties.setProperty(SUBSCRIBE_MODE_NAME, "assign"); return properties; }Parâmetro
Descrição
Método para obter o valor do parâmetro
USER_NAME
Nome de usuário da conta do grupo de consumidores.
AvisoSe não estiver usando o cliente Kafka descrito neste tópico, especifique este parâmetro no seguinte formato:
<Username>-<Consumer group ID>. Exemplo:dtstest-dtsaebpv. Caso contrário, a conexão falhará.No console do DTS, localize a instância de rastreamento de alterações que deseja gerenciar e clique no ID da instância. No painel de navegação à esquerda, clique em Consume Data. Na página exibida, visualize as informações de um grupo de consumidores, como o ID, nome e conta do grupo.
NotaA senha da conta do grupo de consumidores é definida durante a criação do grupo.
PASSWORD_NAME
Senha da conta.
SID_NAME
ID do grupo de consumidores.
GROUP_NAME
Nome do grupo de consumidores. Defina este parâmetro como o ID do grupo de consumidores.
KAFKA_TOPIC
Nome do tópico rastreado da instância de rastreamento de alterações.
No console do DTS, localize a instância de rastreamento de alterações que deseja gerenciar e clique no ID da instância. Na página Basic Information, visualize as informações de tópico e rede. Na seção Basic Information da página de detalhes da tarefa de rastreamento de alterações, obtenha o valor de Subscription Topic. Na seção Network, obtenha o VPC network address (exemplo de formato:
xxx.aliyuncs.com:18003).KAFKA_BROKER_URL_NAME
Endpoint da instância de rastreamento de alterações.
NotaAo rastrear alterações de dados em redes internas, a latência é mínima. Isso se aplica quando a instância do Elastic Compute Service (ECS) onde o cliente Kafka está implantado reside na rede clássica ou na mesma Virtual Private Cloud (VPC) da instância de rastreamento de alterações.
INITIAL_CHECKPOINT_NAME
Checkpoint de consumo dos dados consumidos. O valor é um timestamp UNIX. Exemplo: 1592269238.
Nota-
Salve o checkpoint de consumo pelos seguintes motivos:
-
Se o processo de consumo for interrompido, especifique o checkpoint de consumo no cliente Kafka para retomar o consumo de dados, evitando perda de informações.
-
Ao iniciar o cliente Kafka, especifique o checkpoint de consumo para consumir dados conforme suas necessidades de negócio.
-
-
Se o parâmetro SUBSCRIBE_MODE_NAME estiver definido como subscribe, o parâmetro INITIAL_CHECKPOINT_NAME especificado terá efeito apenas na primeira execução do cliente Kafka.
O checkpoint de consumo dos dados consumidos deve estar dentro do intervalo de dados da instância de rastreamento de alterações. Converta o checkpoint de consumo em um timestamp UNIX. Na lista de tarefas de rastreamento de alterações do DTS, visualize o campo Data Range da tarefa correspondente, que mostra a hora inicial e final dos dados consumíveis. Use esse campo para determinar o intervalo de valores válidos de INITIAL_CHECKPOINT_NAME.
Nota-
Visualize o intervalo de dados da instância de rastreamento de alterações na coluna Data Range na página Change Tracking Tasks.
-
Use um mecanismo de busca para encontrar um conversor de timestamp UNIX.
USE_CONFIG_CHECKPOINT_NAME
Define se o cliente deve ser forçado a consumir dados a partir do checkpoint de consumo especificado. Valor padrão: true. Defina este parâmetro como true para evitar a perda de dados recebidos, mas ainda não processados.
Nenhum
SUBSCRIBE_MODE_NAME
Define se dois ou mais clientes Kafka serão executados para um grupo de consumidores. Para usar esse recurso, defina este parâmetro como subscribe nesses clientes Kafka.
O valor padrão é assign, indicando que o recurso não está em uso. Recomenda-se implantar apenas um cliente Kafka por grupo de consumidores.
Nenhum
-
-
Na barra de menu superior do IntelliJ IDEA, escolha para executar o cliente.
NotaNa primeira execução do IntelliJ IDEA, é necessário algum tempo para carregar e instalar as dependências relevantes.
Resultados no cliente Kafka
O resultado a seguir indica que o cliente Kafka consegue rastrear alterações de dados do banco de dados de origem.
[2020-03-09 10:41:52,408] INFO [Consumer clientId=consumer-1, groupId=dts_xxx] Discovered coordinator xxx (id: xxx rack: null) (org.apache.kafka.clients.consumer.internals.AbstractCoordinator)
[2020-03-09 10:41:57,203] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721711, offset: 1732521, info: 1583721711] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:41:57,571] INFO EtlRecordProcessor: haven't receive records from generator for 5s (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:02,203] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721721, offset: 1732539, info: 1583721721] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:07,204] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721726, offset: 1732544, info: 1583721726] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:12,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721731, offset: 1732548, info: 1583721731] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:17,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721736, offset: 1732554, info: 1583721736] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:22,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721741, offset: 1732559, info: 1583721741] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:27,206] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721746, offset: 1732569, info: 1583721746] (recordprocessor.EtlRecordProcessor)
Remova as barras duplas (//) da string //log.info(ret) na linha 25 do arquivo NotifyDemoDB.java. Em seguida, execute o cliente novamente para visualizar as informações de alteração de dados.
Perguntas frequentes
-
P: Por que preciso registrar o checkpoint de consumo do cliente Kafka?
R: O checkpoint de consumo registrado pelo DTS corresponde ao momento em que o DTS recebe uma operação de commit do cliente Kafka. O checkpoint registrado pode diferir do horário real de consumo. Se uma aplicação de negócio ou o cliente Kafka for interrompido inesperadamente, especifique um checkpoint de consumo preciso para continuar o consumo de dados. Isso evita perda de dados ou consumo duplicado.
Gerenciar o checkpoint de consumo
-
Configure o cliente Kafka para monitorar a troca do módulo de coleta de dados no DTS.
Configure as propriedades do consumidor do cliente Kafka para monitorar a troca do módulo de coleta de dados no DTS. O código a seguir fornece um exemplo de configuração das propriedades do consumidor:
properties.setProperty(ConsumerConfig.INTERCEPTOR_CLASSES_CONFIG, ClusterSwitchListener.class.getName());O código a seguir fornece um exemplo de implementação do ClusterSwitchListener:
public class ClusterSwitchListener implements ClusterResourceListener, ConsumerInterceptor { private final static Logger LOG = LoggerFactory.getLogger(ClusterSwitchListener.class); private ClusterResource originClusterResource = null; private ClusterResource currentClusterResource = null; public ConsumerRecords onConsume(ConsumerRecords records) { return records; } public void close() { } public void onCommit(Map offsets) { } public void onUpdate(ClusterResource clusterResource) { synchronized (this) { originClusterResource = currentClusterResource; currentClusterResource = clusterResource; if (null == originClusterResource) { LOG.info("Cluster updated to " + currentClusterResource.clusterId()); } else { if (originClusterResource.clusterId().equals(currentClusterResource.clusterId())) { LOG.info("Cluster not changed on update:" + clusterResource.clusterId()); } else { LOG.error("Cluster changed"); throw new ClusterSwitchException("Cluster changed from " + originClusterResource.clusterId() + " to " + currentClusterResource.clusterId() + ", consumer require restart"); } } } } public boolean isClusterResourceChanged() { if (null == originClusterResource) { return false; } if (originClusterResource.clusterId().equals(currentClusterResource.clusterId())) { return false; } return true; } public void configure(Map<String, ?> configs) { } public static class ClusterSwitchException extends KafkaException { public ClusterSwitchException(String message, Throwable cause) { super(message, cause); } public ClusterSwitchException(String message) { super(message); } public ClusterSwitchException(Throwable cause) { super(cause); } public ClusterSwitchException() { super(); } } -
Especifique o checkpoint de consumo com base na troca capturada do módulo de coleta de dados no DTS.
Defina o checkpoint inicial de consumo do próximo rastreamento de dados como o timestamp da última entrada de dados rastreada consumida pelo cliente. O código a seguir fornece um exemplo de especificação do checkpoint de consumo:
try{ //do some action } catch (ClusterSwitchListener.ClusterSwitchException e) { reset(); } // Reset the consumption checkpoint. public reset() { long offset = kafkaConsumer.offsetsForTimes(timestamp); kafkaConsumer.seek(tp,offset); }NotaPara mais informações sobre os exemplos, consulte KafkaRecordFetcher.
Mapeamentos entre tipos de dados MySQL e valores dataTypeNumber
Para mais informações, consulte Campo SQL Type.
|
Tipo de dados MySQL |
Valor de dataTypeNumber |
|
MYSQL_TYPE_DECIMAL |
0 |
|
MYSQL_TYPE_INT8 |
1 |
|
MYSQL_TYPE_INT16 |
2 |
|
MYSQL_TYPE_INT32 |
3 |
|
MYSQL_TYPE_FLOAT |
4 |
|
MYSQL_TYPE_DOUBLE |
5 |
|
MYSQL_TYPE_NULL |
6 |
|
MYSQL_TYPE_TIMESTAMP |
7 |
|
MYSQL_TYPE_INT64 |
8 |
|
MYSQL_TYPE_INT24 |
9 |
|
MYSQL_TYPE_DATE |
10 |
|
MYSQL_TYPE_TIME |
11 |
|
MYSQL_TYPE_DATETIME |
12 |
|
MYSQL_TYPE_YEAR |
13 |
|
MYSQL_TYPE_DATE_NEW |
14 |
|
MYSQL_TYPE_VARCHAR |
15 |
|
MYSQL_TYPE_BIT |
16 |
|
MYSQL_TYPE_TIMESTAMP_NEW |
17 |
|
MYSQL_TYPE_DATETIME_NEW |
18 |
|
MYSQL_TYPE_TIME_NEW |
19 |
|
MYSQL_TYPE_JSON |
245 |
|
MYSQL_TYPE_DECIMAL_NEW |
246 |
|
MYSQL_TYPE_ENUM |
247 |
|
MYSQL_TYPE_SET |
248 |
|
MYSQL_TYPE_TINY_BLOB |
249 |
|
MYSQL_TYPE_MEDIUM_BLOB |
250 |
|
MYSQL_TYPE_LONG_BLOB |
251 |
|
MYSQL_TYPE_BLOB |
252 |
|
MYSQL_TYPE_VAR_STRING |
253 |
|
MYSQL_TYPE_STRING |
254 |
|
MYSQL_TYPE_GEOMETRY |
255 |
Mapeamentos entre tipos de dados Oracle e valores dataTypeNumber
|
Tipo de dados Oracle |
Valor de dataTypeNumber |
|
VARCHAR2/NVARCHAR2 |
1 |
|
NUMBER/FLOAT |
2 |
|
LONG |
8 |
|
DATE |
12 |
|
RAW |
23 |
|
LONG_RAW |
24 |
|
UNDEFINED |
29 |
|
XMLTYPE |
58 |
|
ROWID |
69 |
|
CHAR e NCHAR |
96 |
|
BINARY_FLOAT |
100 |
|
BINARY_DOUBLE |
101 |
|
CLOB/NCLOB |
112 |
|
BLOB |
113 |
|
BFILE |
114 |
|
TIMESTAMP |
180 |
|
TIMESTAMP_WITH_TIME_ZONE |
181 |
|
INTERVAL_YEAR_TO_MONTH |
182 |
|
INTERVAL_DAY_TO_SECOND |
183 |
|
UROWID |
208 |
|
TIMESTAMP_WITH_LOCAL_TIME_ZONE |
231 |
Mapeamentos entre tipos de dados PostgreSQL e valores dataTypeNumber
|
Tipo de dados PostgreSQL |
Valor de dataTypeNumber |
|
INT2/SMALLINT |
21 |
|
INT4/INTEGER/SERIAL |
23 |
|
INT8/BIGINT |
20 |
|
CHARACTER |
18 |
|
CHARACTER VARYING |
1043 |
|
REAL |
700 |
|
DOUBLE PRECISION |
701 |
|
NUMERIC |
1700 |
|
MONEY |
790 |
|
DATE |
1082 |
|
TIME/TIME WITHOUT TIME ZONE |
1083 |
|
TIME WITH TIME ZONE |
1266 |
|
TIMESTAMP/TIMESTAMP WITHOUT TIME ZONE |
1114 |
|
TIMESTAMP WITH TIME ZONE |
1184 |
|
BYTEA |
17 |
|
TEXT |
25 |
|
JSON |
114 |
|
JSONB |
3082 |
|
XML |
142 |
|
UUID |
2950 |
|
POINT |
600 |
|
LSEG |
601 |
|
PATH |
602 |
|
BOX |
603 |
|
POLYGON |
604 |
|
LINE |
628 |
|
CIDR |
650 |
|
CIRCLE |
718 |
|
MACADDR |
829 |
|
INET |
869 |
|
INTERVAL |
1186 |
|
TXID_SNAPSHOT |
2970 |
|
PG_LSN |
3220 |
|
TSVECTOR |
3614 |
|
TSQUERY |
3615 |