Use o conector JDBC SQL do Apache Flink para gravar dados em streaming ou em lote em um cluster do AnalyticDB for MySQL Data Warehouse Edition. Este guia abrange o Flink 1,11 e versões posteriores.
Este guia usa o Apache Flink SQL. Para gravar dados com a API Flink Java Database Connectivity (JDBC), consulte JDBC Connector.
Pré-requisitos
Antes de começar, verifique se você tem:
Um cluster do AnalyticDB for MySQL Data Warehouse Edition com banco de dados e tabela prontos para receber dados. Consulte CREATE DATABASE e CREATE TABLE
Um cluster Apache Flink na versão 1,11 ou superior
O JAR do conector JDBC do Flink e um JAR de driver MySQL (versão 5.1.40 ou superior) implantados no diretório
${Flink deployment directory}/libem todos os nós do Flink, com o cluster reiniciado após a implantação(Apenas clusters em modo elástico) ENI ativada na página Cluster Information, em Network Information
Ativar ou desativar a ENI interrompe todas as conexões com o banco de dados por aproximadamente 2 minutos. Durante esse período, as operações de leitura e gravação ficam indisponíveis.
Baixe os arquivos JAR necessários
Conector JDBC do Flink
Baixe o JAR do conector correspondente à sua versão do Flink. Para versões não listadas abaixo, consulte JDBC SQL Connector.
|
Versão do Flink |
Arquivo JAR |
|
1,11 |
|
|
1,12 |
|
|
1,13 |
Driver MySQL
Baixe um driver MySQL (versão 5.1.40 ou superior) em mysql/mysql-connector-java.
Após colocar ambos os JARs no diretório lib de cada nó do Flink, reinicie o cluster. Para instruções de inicialização do cluster, consulte Step 2: Start a cluster.
Compatibilidade de versões anteriores do Flink
Para versões do Flink anteriores à 1,11:
Flink 1,9 e 1,10: Flink v1.10 JDBC connector
Flink 1,8 e anteriores: Flink v1.8 JDBCAppendTableSink
Etapa 1: Preparar os dados de origem
Este exemplo usa um arquivo CSV local como source de dados.
-
Em qualquer nó do Flink, crie o arquivo
/root/data.csvcom o seguinte conteúdo:0,json00,20 1,json01,21 2,json02,22 3,json03,23 4,json04,24 5,json05,25 6,json06,26 7,json07,27 8,json08,28 9,json09,29 Copie o arquivo para
/root/data.csvem todos os outros nós do Flink. O caminho deve ser idêntico em cada nó.
Etapa 2: Gravar dados no AnalyticDB for MySQL
Execute todos os comandos SQL desta etapa na CLI do Flink SQL Client. Para iniciá-la, consulte Starting the SQL Client CLI.
Crie a tabela de origem
Crie uma tabela de origem no Flink para ler o arquivo CSV:
CREATE TABLE IF NOT EXISTS csv_person (
`user_id` STRING,
`user_name` STRING,
`age` INT
) WITH (
'connector' = 'filesystem',
'path' = 'file:///root/data.csv',
'format' = 'csv',
'csv.ignore-parse-errors' = 'true',
'csv.allow-comments' = 'true'
);
Os nomes das colunas e os tipos de dados devem corresponder aos da tabela de destino no AnalyticDB for MySQL. O valor de path deve apontar para o mesmo caminho absoluto em todos os nós do Flink. Para outras opções de conector, consulte FileSystem SQL Connector.
Crie a tabela de resultado
Crie uma tabela de resultado no Flink mapeada para a tabela de destino no AnalyticDB for MySQL:
CREATE TABLE mysql_person (
user_id STRING,
user_name STRING,
age INT
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://<endpoint>:<port>/<db_name>?useServerPrepStmts=false&rewriteBatchedStatements=true',
'table-name' = '<table_name>',
'username' = '<username>',
'password' = '<password>',
'sink.buffer-flush.max-rows' = '10',
'sink.buffer-flush.interval' = '1s'
);
Substitua os placeholders antes de executar esta instrução:
|
Placeholder |
Descrição |
Exemplo |
|
|
Endpoint do cluster. Para usar um endpoint público, primeiro solicite um. |
|
|
|
Porta do cluster |
|
|
|
Nome do banco de dados de destino |
|
|
|
Nome da tabela de destino |
|
|
|
Conta de banco de dados com permissões de gravação. Execute SHOW GRANTS para verificar as permissões ou GRANT para atribuí-las. |
— |
|
|
Senha da conta do banco de dados |
— |
Parâmetros do conector
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Sim |
— |
Defina como |
|
|
Sim |
— |
URL JDBC do cluster. A string de consulta |
|
|
Sim |
— |
Nome da tabela de destino. |
|
|
Sim |
— |
Conta de banco de dados com permissões de gravação. |
|
|
Sim |
— |
Senha da conta. |
|
|
Não |
— |
Número máximo de linhas por lote. Defina como |
|
|
Não |
— |
Tempo máximo entre liberações. Defina como |
Para todas as opções disponíveis do conector, consulte Connector options.
Comportamento de liberação (flush)
Quando sink.buffer-flush.max-rows e sink.buffer-flush.interval têm valores diferentes de zero, a liberação ocorre assim que uma das condições é atendida:
O número de linhas no buffer atinge
sink.buffer-flush.max-rowsO tempo decorrido desde a última liberação atinge
sink.buffer-flush.interval
Execute a importação
Execute a seguinte instrução para iniciar a gravação de dados da tabela de origem no AnalyticDB for MySQL:
INSERT INTO mysql_person SELECT user_id, user_name, age FROM csv_person;
Se a tabela de destino tiver uma chave primária e os dados de origem contiverem valores duplicados dessa chave, INSERT INTO se comportará como INSERT IGNORE INTO — as linhas duplicadas são ignoradas em vez de inseridas novamente. Consulte INSERT INTO.
Etapa 3: Verifique os dados
Após a conclusão do job, faça login no banco de dados tpch do seu cluster AnalyticDB for MySQL e execute:
SELECT * FROM person;