Todos os produtos
Search
Central de documentação

AnalyticDB:Importar dados do Apache Flink

Última atualização: Jun 27, 2026

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.

Nota

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}/lib em 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

Importante

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

flink-connector-jdbc_2.11-1.11.0.jar

1,12

flink-connector-jdbc_2.11-1.12.0.jar

1,13

flink-connector-jdbc_2.11-1.13.0.jar

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:

Etapa 1: Preparar os dados de origem

Este exemplo usa um arquivo CSV local como source de dados.

  1. Em qualquer nó do Flink, crie o arquivo /root/data.csv com 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
  2. Copie o arquivo para /root/data.csv em 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'
);
Nota

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>

Endpoint do cluster. Para usar um endpoint público, primeiro solicite um.

am-**********.ads.aliyuncs.com

<port>

Porta do cluster

3306

<db_name>

Nome do banco de dados de destino

tpch

<table_name>

Nome da tabela de destino

person

<username>

Conta de banco de dados com permissões de gravação. Execute SHOW GRANTS para verificar as permissões ou GRANT para atribuí-las.

<password>

Senha da conta do banco de dados

Parâmetros do conector

Parâmetro

Obrigatório

Padrão

Descrição

connector

Sim

Defina como jdbc.

url

Sim

URL JDBC do cluster. A string de consulta useServerPrepStmts=false&rewriteBatchedStatements=true habilita gravações em lote para melhorar o throughput e reduzir a carga do cluster.

table-name

Sim

Nome da tabela de destino.

username

Sim

Conta de banco de dados com permissões de gravação.

password

Sim

Senha da conta.

sink.buffer-flush.max-rows

Não

Número máximo de linhas por lote. Defina como 0 para liberar apenas com base no intervalo de tempo. Exemplos diferentes de zero: 1000, 2000. Evite definir como 0, pois isso degrada o throughput de gravação durante consultas concorrentes.

sink.buffer-flush.interval

Não

Tempo máximo entre liberações. Defina como 0 para liberar apenas quando max-rows for atingido. Exemplos diferentes de zero: 1d, 1h, 1min, 1s, 1ms. Evite definir como 0, pois os dados podem sofrer atraso em períodos de baixo tráfego.

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

  • O 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;
Nota

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;