Todos os produtos
Search
Central de documentação

Data Transmission Service:Consumir dados com o flink-dts-connector

Última atualização: Jun 27, 2026

Após configurar uma instância de assinatura, use o flink-dts-connector para consumir dados dessa instância com um cliente Flink. Este tópico descreve como usar o flink-dts-connector.

Observações de uso

  • Este conector é compatível com clientes Flink que usam a DataStream API, a Table API ou SQL.

  • Se o cliente Flink usar a Table API ou SQL, será possível consumir dados de apenas uma tabela por configuração. Para consumir dados de várias tabelas, configure e execute múltiplas tarefas independentes.

Procedimento

Este tópico usa o IntelliJ IDEA Community Edition 2020.1 para Windows como exemplo para demonstrar o uso do flink-dts-connector no consumo de dados de uma instância de assinatura.

  1. Crie uma tarefa de rastreamento de alterações. Para mais informações, consulte os tópicos relevantes em Visão geral do rastreamento de alterações.

  2. Crie um ou mais grupos de consumidores. Para mais informações, consulte Adicionar um grupo de consumidores.

  3. Baixe e descompacte o flink-dts-connector.

  4. Abra o IntelliJ IDEA e clique em Open or Import.

    打开工程

  5. Na caixa de diálogo exibida, acesse o diretório descompactado do flink-dts-connector, expanda as pastas e localize o arquivo Project Object Model (pom.xml).

    pom模型

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

  7. Adicione a seguinte dependência ao arquivo pom.xml:

    <dependency>
          <groupId>com.alibaba.flink</groupId>
          <artifactId>flink-dts-connector</artifactId>
          <version>1.1.1-SNAPSHOT</version>
          <classifier>jar-with-dependencies</classifier>
    </dependency>
  8. No IntelliJ IDEA, expanda as pastas e selecione o arquivo Java correspondente ao tipo de programa do seu conector Flink.

    • Se o cliente Flink usar a DataStream API, clique duas vezes no arquivo flink-dts-connector-master\src\test\java\com\alibaba\flink\connectors\dts\datastream\DtsExample.java e execute as etapas a seguir:

      1. Na parte superior da janela do IntelliJ IDEA, clique em Run.run图标

      2. No menu exibido, clique em DtsExample > Edit.edit

      3. No campo Program arguments, insira os parâmetros e seus valores conforme o exemplo abaixo. Em seguida, clique em Run para iniciar o flink-dts-connector.

        Nota

        Para obter mais informações sobre os parâmetros e como obter seus valores, consulte Parâmetros.

        --broker-url dts-cn-******.******.***:****** --topic cn_hangzhou_rm_**********_dtstest_version2 --sid dts****** --user dtstest --password Test123456 --checkpoint 1624440043
      4. A saída confirma que o cliente consegue assinar alterações de dados do banco de dados de origem.数据变更信息(DataStream)

        Nota

        Para consultar registros específicos de alteração de dados, faça login na interface do Task Manager do seu cliente Flink.

    • Se o cliente Flink usar a Table API ou SQL, clique duas vezes no arquivo flink-dts-connector-master\src\test\java\com\alibaba\flink\connectors\dts\sql\DtsTableISelectTCaseTest.java e execute as etapas a seguir:

      Nota

      Cada arquivo DtsTableISelectTCaseTest.java pode ser configurado para consumir dados de apenas uma tabela. Para consumir dados de várias tabelas, repita a configuração e execute múltiplas tarefas independentes.

      1. Adicione duas barras (//) para comentar a linha de código especificada.注释掉一行

      2. Configure as informações da tabela única da qual deseja consumir dados. Há suporte para instruções SQL.

      3. Defina os parâmetros da instância de assinatura. Para obter mais informações sobre os parâmetros e como obter seus valores, consulte Parâmetros.table api和sql的参数配置

      4. Na parte superior da janela do IntelliJ IDEA, clique em Run'DtsTableISelectTCaseTest' para iniciar o flink-dts-connector.

      5. A saída confirma que o cliente consegue assinar alterações de dados do banco de dados de origem.tableapi和sql-数据变更信息

        Nota

        Para consultar registros específicos de alteração de dados, faça login na interface do Task Manager do seu cliente Flink.

Parâmetros

DtsExample.java

DtsTableISelectTCaseTest.java

Descrição

Método de consulta

broker-url

dts.server

O endpoint de rede e a porta da instância de assinatura.

Nota
  • Se a instância ECS onde o cliente Flink está implantado e a instância de assinatura estiverem na mesma rede clássica ou Virtual Private Cloud (VPC), recomenda-se o uso de um endpoint interno para a assinatura de dados, visando minimizar a latência de rede.

  • Não recomendamos o uso de endpoints públicos devido à possível instabilidade da rede.

No console DTS, clique em ID da instância de assinatura desejada. Na página Basic Information, obtenha as informações de Topic e Network.

topic

topic

O tópico da instância de assinatura.

sid

dts.sid

O ID do grupo de consumidores.

No console DTS, clique em ID da instância de assinatura desejada. No painel de navegação à esquerda, clique em Consume Data. É possível obter o Consumer Group ID/Name e a Account do grupo de consumidores.

Nota

A senha da conta do grupo de consumidores é definida durante a criação do grupo.

user

dts.user

A conta do grupo de consumidores.

Aviso

Se você não estiver usando o flink-dts-connector fornecido neste tópico, defina o nome de usuário no formato <consumer_group_account>-<consumer_group_ID>, como dtstest-dtsaebpv. Caso contrário, a conexão falhará.

password

dts.password

A senha da conta.

checkpoint

dts.checkpoint

O offset do consumidor. Este parâmetro especifica o timestamp Unix do primeiro registro de dados a ser consumido. Exemplo: 1624440043.

Nota

Use o offset do consumidor para:

  • Retomar o consumo a partir do último offset após uma interrupção do aplicativo, evitando perda de dados.

  • Definir um offset para iniciar o consumo de dados a partir de um ponto específico no tempo.

O horário inicial para consumo deve estar dentro do intervalo de dados da instância de assinatura e precisa ser convertido para um timestamp Unix.

Nota
  • Visualize o intervalo de dados da instância de assinatura na coluna Data Range da lista de tarefas de assinatura.

  • Use um mecanismo de busca para encontrar uma ferramenta de conversão de timestamp Unix.

N/A

dts-cdc.table.name

O objeto assinado. Especifique apenas uma única tabela. Os requisitos de formato são:

  • Se o Database Type for MySQL, PolarDB for MySQL, PolarDB-X 1.0 ou PolarDB-X 2.0, use o formato <database_name>.<table_name>.

  • Para outros tipos de banco de dados, use o formato <schema_name>.<table_name>.

No console DTS, clique em ID da instância de assinatura desejada. Na parte superior da página Basic Information ou Task Management, clique em View Objects para localizar o banco de dados e a tabela assinados.

Perguntas frequentes

Mensagem de erro

Possível causa

Solução

Cluster changed from *** to ***, consumer require restart.

Um failover no módulo DStore, usado pelo Data Transmission Service (DTS) para ler dados incrementais, causa a perda do offset do consumidor no cliente Flink.

Não reinicie o cliente. Em vez disso, consulte o offset do consumidor e passe o valor novamente para o parâmetro checkpoint ou dts.checkpoint no arquivo DtsExample.java ou DtsTableISelectTCaseTest.java para retomar o consumo de dados.