Todos os produtos
Search
Central de documentação

Data Transmission Service:Usar um cliente Kafka para consumir dados rastreados

Última atualização: Aug 24, 2026

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.

    Nota

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

    Aviso

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

Nota
  • Clique em code 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>
Tabela 1 Descrição do processo

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
  • Se a instância de origem for um banco de dados Oracle autogerenciado, ative o log suplementar para todas as colunas. Isso garante que o cliente consuma os dados rastreados com êxito e assegura a integridade da pré-imagem e da pós-imagem.

  • Se a instância de origem não for um banco de dados Oracle autogerenciado, o DTS não garante a integridade da pré-imagem. Valide a pré-imagem obtida.

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.

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.

  1. Crie uma instância de rastreamento de alterações. Para mais informações, consulte Change tracking overview.

  2. Crie um ou mais grupos de consumidores. Para mais informações, consulte Add a consumer group.

  3. Baixe o pacote da demonstração do cliente Kafka e descompacte-o.

    Nota

    Clique em code e selecione Download ZIP para baixar o pacote.

  4. Abra o IntelliJ IDEA. Na janela exibida, clique em Open.

  5. 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, selecione pom.xml e clique em OK.

  6. Na caixa de diálogo exibida, selecione Open as Project.

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

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

    Aviso

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

    Nota

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

    Nota

    Ao 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

  9. Na barra de menu superior do IntelliJ IDEA, escolha Run > Run para executar o cliente.

    Nota

    Na 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

  1. 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();
            }
    
        }
  2. 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);
    }
    Nota

    Para 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