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.
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.
Crie um ou mais grupos de consumidores. Para mais informações, consulte Adicionar um grupo de consumidores.
Baixe e descompacte o flink-dts-connector.
-
Abra o IntelliJ IDEA e clique em Open or Import.

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

Na caixa de diálogo exibida, selecione Open as Project.
-
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> -
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:
Na parte superior da janela do IntelliJ IDEA, clique em Run.

No menu exibido, clique em .

-
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.
NotaPara 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 -
A saída confirma que o cliente consegue assinar alterações de dados do banco de dados de origem.
NotaPara 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:
NotaCada arquivo
DtsTableISelectTCaseTest.javapode 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.Adicione duas barras (
//) para comentar a linha de código especificada.
Configure as informações da tabela única da qual deseja consumir dados. Há suporte para instruções SQL.
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.

Na parte superior da janela do IntelliJ IDEA, clique em Run'DtsTableISelectTCaseTest' para iniciar o flink-dts-connector.
-
A saída confirma que o cliente consegue assinar alterações de dados do banco de dados de origem.
NotaPara 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 |
|
|
|
O endpoint de rede e a porta da instância de assinatura. Nota
|
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. |
|
|
|
O tópico da instância de assinatura. |
|
|
|
|
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. |
|
|
|
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 |
|
|
|
|
A senha da conta. |
|
|
|
|
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:
|
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
|
|
N/A |
|
O objeto assinado. Especifique apenas uma única tabela. Os requisitos de formato são:
|
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 |
|
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 |