Use o flink-dts-connector para criar um programa Flink que consuma dados rastreados de uma instância de rastreamento de alterações do Data Transmission Service (DTS).
Observações de uso
O conector oferece suporte a programas Flink que usam a DataStream API ou a Table API e SQL.
Se o programa Flink usar a Table API e SQL, configure-o para consumir dados de apenas uma tabela por vez. Para consumir dados de várias tabelas, configure e execute uma tarefa separada para cada uma.
Procedimento
O exemplo a seguir usa o IntelliJ IDEA Community Edition 2020.1 para Windows para demonstrar o consumo de dados rastreados com o flink-dts-connector.
Crie uma instância de rastreamento de alterações. Para mais informações, consulte Criar uma instância de rastreamento de alterações para uma instância ApsaraDB RDS for MySQL, Criar uma instância de rastreamento de alterações para um cluster PolarDB for MySQL ou Criar uma instância de rastreamento de alterações para uma instância ApsaraDB for Oracle.
Crie um ou mais grupos de consumidores. Para mais informações, consulte Criar um grupo de consumidores.
Baixe o arquivo flink-dts-connector e descompacte-o.
Inicie o IntelliJ IDEA e clique em Open or Import.
Na caixa de diálogo exibida, acesse o diretório onde você descompactou o arquivo flink-dts-connector. Expanda as pastas para localizar o arquivo Project Object Model (POM): 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 do projeto e selecione o arquivo Java apropriado conforme o tipo de API Flink usada pelo programa.
-
Se o programa 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 barra de menus superior, escolha Run > Run....
Na janela pop-up, clique em .
-
No campo Program arguments, insira os parâmetros e seus valores conforme o exemplo a seguir. Em seguida, clique em Run para iniciar o flink-dts-connector.
NotaPara obter detalhes sobre os parâmetros e como encontrar 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 indica que o programa está rastreando com êxito as alterações de dados do banco de dados source. Após a inicialização, os registros de alteração de dados no DataStream serão semelhantes ao exemplo a seguir.
LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006303@289540@2688@1625045211000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006307@290047@2688@1625045212000]} LazyParseRecord {operationType [UPDATE], checkpoint [0@12006305@290016@2688@1625045212000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006308@290047@2688@1625045214000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006309@290047@2688@1625045215000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006310@290047@2688@1625045216000]} LazyParseRecord {operationType [HEARTBEAT], checkpoint [0@12006311@290047@2688@1625045217000]}NotaPara visualizar os detalhes dos registros de alteração de dados, faça login na interface do Task Manager do seu programa Flink.
-
Se o programa Flink usar a Table API e 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:
NotaUm único arquivo
DtsTableISelectTCaseTest.javapermite configurar e consumir dados rastreados de apenas uma tabela. Para consumir dados de várias tabelas, configure e execute uma tarefa separada para cada uma.-
Comente a linha
properties.loadadicionando o prefixo//./*Loads the parameter values that you set on the platform page into the Properties object.*/ //properties.load(new StringReader(new String(Files.readAllBytes(Paths.get(configFilePath)), StandardCharsets.UTF_8))); Defina a única tabela da qual deseja consumir dados.
-
Configure os parâmetros da instância de rastreamento de alterações. Para obter detalhes sobre os parâmetros e como encontrar seus valores, consulte Parâmetros.
public static void main(String[] args) throws Exception { setup(args); final String createTable = "create table `dts` (\n" + " `ts` TIMESTAMP(3) METADATA FROM 'timestamp',\n" + " `id` bigint,\n" + " `name` varchar,\n" + " `age` bigint,\n" + " WATERMARK FOR ts AS ts - INTERVAL '5' SECOND" + ") with (\n" + "'connector' = 'dts'," + "'dts.server' = 'dts-cn-xxx:18001'," + "'topic' = 'cn_hangzhou_rm_xxx_dtstest_version2'," + "'dts.sid' = 'dtsxxx', " + "'dts.user' = 'dtstest', " + "'dts.password' = 'xxx'," + "'dts.checkpoint' = '1624440043', " + "'dts-cdc.table.name' = 'dtstestdata.order'," + "'format' = 'dts-cdc')"; } Na parte superior da interface do IntelliJ IDEA, clique em Run'DtsTableISelectTCaseTest' para iniciar o flink-dts-connector.
-
A saída indica que o programa está rastreando com êxito as alterações de dados do banco de dados source. Após a inicialização, o terminal exibe registros de changelog. Os eventos de atualização aparecem em pares: o valor antes da atualização (-U) e o valor após a atualização (+U).
######> (false,-U(2021-06-23T20:32:17.391,null,null,null)) ######> (true,+U(2021-06-23T20:32:17.391,null,null,null)) ######> (false,-U(2021-06-23T20:32:45.604,null,null,null)) ######> (true,+U(2021-06-23T20:32:45.604,null,null,null)) ######> (false,-U(2021-06-30T17:26:52.201,null,null,null)) ######> (true,+U(2021-06-30T17:26:52.201,null,null,null)) ######> (false,-U(2021-06-30T19:19:26.975,null,null,null)) ######> (true,+U(2021-06-30T19:19:26.975,null,null,null))NotaPara visualizar os detalhes dos registros de alteração de dados, faça login na interface do Task Manager do seu programa Flink.
-
-
Parâmetros
|
Parâmetro da DataStream API |
Parâmetro da Table API |
Descrição |
Como obter |
|
|
|
Endpoint e porta da instância de rastreamento de alterações. Nota
|
No console do Data Transmission Service (DTS), clique em ID da instância de rastreamento de alterações desejada. Na página View Task Settings, localize o Topic, o endpoint e a porta. |
|
|
|
Tópico rastreado da instância de rastreamento de alterações. |
|
|
|
|
ID do grupo de consumidores. |
No console do DTS, clique em ID da instância de rastreamento de alterações desejada e, em seguida, clique em Consume Data. Localize o ID e a Account do grupo de consumidores. Nota
A senha do nome de usuário do grupo de consumidores é aquela definida durante a criação do grupo. |
|
|
|
Nome de usuário do grupo de consumidores. Aviso
O flink-dts-connector fornecido neste tópico lida automaticamente com a formatação necessária do nome de usuário. Se você usar um cliente diferente, formate manualmente o nome de usuário como |
|
|
|
|
Senha do nome de usuário. |
|
|
|
|
Timestamp Unix a partir do qual o flink-dts-connector começa a consumir dados, por exemplo, 1624440043. Nota
O offset do consumidor é útil nos seguintes cenários:
|
O offset do consumidor deve estar dentro do intervalo de dados da instância de rastreamento de alterações. Na página View Task Settings da instância de rastreamento de alterações do DTS, localize os horários de início e fim no campo Data Range. Defina o offset do consumidor para um horário dentro desse intervalo e converta-o para um timestamp Unix. Nota
Use um mecanismo de busca online para encontrar um conversor de timestamp Unix. |
|
N/A |
|
Tabela cujas alterações devem ser rastreadas. Apenas uma tabela é suportada por tarefa. Use o seguinte formato:
|
No console do DTS, clique em ID da instância de rastreamento de alterações desejada. Na página View Task Settings, clique em View Objects no canto superior direito para localizar os nomes do banco de dados e das tabelas dos objetos. |
Perguntas frequentes
|
Mensagem de erro |
Possível causa |
Solução |
|
O módulo DStore, usado pelo DTS para ler dados incrementais, passou por um failover. Isso resulta na perda do offset do consumidor do programa Flink. |
Não reinicie o programa. Em vez disso, localize o último offset de consumidor conhecido e passe-o como o parâmetro |