Use o conector Flink do MaxCompute para gravar dados do Flink em tabelas padrão e delta tables no MaxCompute e simplificar a ingestão de dados. Este tópico descreve os recursos do conector e detalha os procedimentos de gravação.
Contexto
-
Modos de gravação compatíveis
O conector Flink oferece dois modos de gravação:
upserteinsert. No modoupsert, agrupe os fluxos de dados das seguintes formas:Agrupamento por chave primária
-
Agrupamento por campo de partição
Embora adequado para um grande número de partições, o agrupamento por campo de partição pode causar skew de dados.
Para obter detalhes sobre o procedimento de gravação
upserte os parâmetros recomendados, consulte Ingestão de dados em tempo real em um data warehouse.Especifique o modo de gravação por meio dos parâmetros do conector Flink. Para a lista completa de parâmetros, consulte Apêndice: Parâmetros do conector Flink.
Defina o intervalo de checkpoint para jobs de gravação upsert do Flink em pelo menos 3 minutos. Intervalos menores reduzem a eficiência da gravação e geram muitos arquivos pequenos.
-
A tabela a seguir apresenta o mapeamento de tipos de dados entre Realtime Compute for Apache Flink e MaxCompute.
Tipo de dado Flink
Tipo de dado MaxCompute
CHAR(p)
CHAR(p)
VARCHAR(p)
VARCHAR(p)
STRING
STRING
BOOLEAN
BOOLEAN
TINYINT
TINYINT
SMALLINT
SMALLINT
INT
INT
BIGINT
BIGINT
FLOAT
FLOAT
DOUBLE
DOUBLE
DECIMAL(p, s)
DECIMAL(p, s)
DATE
DATE
TIMESTAMP(9) WITHOUT TIME ZONE, TIMESTAMP_LTZ(9)
TIMESTAMP
TIMESTAMP(3) WITHOUT TIME ZONE, TIMESTAMP_LTZ(3)
DATETIME
BYTES
BINARY
ARRAY<T>
ARRAY<T>
MAP<K, V>
MAP<K, V>
ROW
STRUCT
NotaO tipo de dado TIMESTAMP do Flink não inclui informações de fuso horário, enquanto o TIMESTAMP do MaxCompute inclui. Essa diferença pode causar uma discrepância de 8 horas. Para alinhar os timestamps, use TIMESTAMP_LTZ(9).
-- Flink SQL CREATE TEMPORARY TABLE odps_source( id BIGINT NOT NULL COMMENT 'ID', created_time TIMESTAMP NOT NULL COMMENT 'Creation time', updated_time TIMESTAMP_LTZ(9) NOT NULL COMMENT 'Update time', PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'maxcompute', ... );
Gravar dados de um cluster Flink autogerenciado
-
Pré-requisitos: Crie uma tabela no MaxCompute.
Crie primeiro uma tabela no MaxCompute para receber os dados do Flink. O exemplo abaixo demonstra esse processo com a criação de duas tabelas (uma delta table não particionada e uma tabela particionada). Para mais informações sobre as configurações de propriedades da tabela, consulte parâmetros de delta table.
-- Create a non-partitioned delta table. CREATE TABLE mf_flink_tt ( id BIGINT not null, name STRING, age INT, status BOOLEAN, primary key (id) ) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; --Create a partitioned delta table. CREATE TABLE mf_flink_tt_part ( id BIGINT not null, name STRING, age INT, status BOOLEAN, primary key (id) ) partitioned by (dd string, hh string) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; -
Configure um cluster Flink open source. O conector é compatível com as versões 1,13, 1,15, 1,16 e 1,17 do Flink. Selecione o conector Flink correspondente à sua versão:
NotaO conector Flink para Flink 1.16 é compatível com Flink 1.17.
Este tópico usa o conector Flink para Flink 1.13 como exemplo. Baixe e descompacte o pacote.
-
Baixe o conector Flink e adicione-o ao pacote do seu cluster Flink.
Baixe o pacote JAR do conector Flink para o seu ambiente local.
-
Adicione o pacote JAR do conector Flink ao diretório lib do pacote de instalação descompactado do Flink.
mv flink-connector-odps-1.13-shaded.jar $FLINK_HOME/lib/flink-connector-odps-1.13-shaded.jar
-
Inicie o serviço da instância Flink.
cd $FLINK_HOME/bin ./start-cluster.sh -
Inicie o cliente Flink.
cd $FLINK_HOME/bin ./sql-client.sh -
Crie uma tabela Flink e configure os parâmetros do conector.
Crie uma tabela Flink e configure seus parâmetros usando Flink SQL ou a API DataStream. As seções a seguir fornecem exemplos essenciais para ambas as abordagens.
Flink SQL
-
No editor Flink SQL, execute os comandos abaixo para criar uma tabela e configurar os parâmetros.
-- Register a non-partitioned table in Flink SQL. CREATE TABLE mf_flink ( id BIGINT, name STRING, age INT, status BOOLEAN, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_tt', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj' ); -- Register a partitioned table in Flink SQL. CREATE TABLE mf_flink_part ( id BIGINT, name STRING, age INT, status BOOLEAN, dd STRING, hh STRING, PRIMARY KEY(id) NOT ENFORCED ) PARTITIONED BY (`dd`,`hh`) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_tt_part', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj' ); -
Grave dados na tabela Flink e consulte a tabela MaxCompute para verificar o resultado.
-- Insert data into the non-partitioned table in the Flink SQL client. INSERT INTO mf_flink VALUES (1,'Danny',27, false); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt; +------------+------+------+--------+ | id | name | age | status | +------------+------+------+--------+ | 1 | Danny | 27 | false | +------------+------+------+--------+ -- Insert data into the non-partitioned table in the Flink SQL client to update the record. INSERT INTO mf_flink VALUES (1,'Danny',28, false); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt; +------------+------+------+--------+ | id | name | age | status | +------------+------+------+--------+ | 1 | Danny | 28 | false | +------------+------+------+--------+ -- Insert data into the partitioned table in the Flink SQL client. INSERT INTO mf_flink_part VALUES (1,'Danny',27, false, '01','01'); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01; +------------+------+------+--------+----+----+ | id | name | age | status | dd | hh | +------------+------+------+--------+----+----+ | 1 | Danny | 27 | false | 01 | 01 | +------------+------+------+--------+----+----+ -- Insert data into the partitioned table in the Flink SQL client to update the record. INSERT INTO mf_flink_part VALUES (1,'Danny',30, false, '01','01'); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01; +------------+------+------+--------+----+----+ | id | name | age | status | dd | hh | +------------+------+------+--------+----+----+ | 1 | Danny | 30 | false | 01 | 01 | +------------+------+------+--------+----+----+
DataStream API
-
Para usar a API DataStream, adicione a seguinte dependência.
<dependency> <groupId>com.aliyun.odps</groupId> <artifactId>flink-connector-maxcompute</artifactId> <version>xxx</version> <scope>system</scope> <systemPath>${mvn_project.basedir}/lib/flink-connector-maxcompute-xxx-shaded.jar</systemPath> </dependency>NotaSubstitua "xxx" pelo número real da versão.
-
O código de exemplo a seguir mostra como criar uma tabela e configurar parâmetros.
package com.aliyun.odps.flink.examples; import org.apache.flink.configuration.Configuration; import org.apache.flink.odps.table.OdpsOptions; import org.apache.flink.odps.util.OdpsConf; import org.apache.flink.odps.util.OdpsPipeline; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.Table; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; import org.apache.flink.table.data.RowData; public class Examples { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(120 * 1000); StreamTableEnvironment streamTableEnvironment = StreamTableEnvironment.create(env); Table source = streamTableEnvironment.sqlQuery("SELECT * FROM source_table"); DataStream<RowData> input = streamTableEnvironment.toAppendStream(source, RowData.class); Configuration config = new Configuration(); config.set(OdpsOptions.SINK_OPERATION, "upsert"); config.set(OdpsOptions.UPSERT_COMMIT_THREAD_NUM, 8); config.set(OdpsOptions.UPSERT_MAJOR_COMPACT_MIN_COMMITS, 100); OdpsConf odpsConfig = new OdpsConf("accessid", "accesskey", "endpoint", "project", "tunnel endpoint"); OdpsPipeline.Builder builder = OdpsPipeline.builder(); builder.projectName("sql2_isolation_2a") .tableName("user_ledger_portfolio") .partition("") .configuration(config) .odpsConf(odpsConfig) .sink(input, false); env.execute(); } }
-
Gravar dados do Flink totalmente gerenciado
-
Pré-requisitos: Crie uma tabela no MaxCompute.
Crie uma tabela de destino no MaxCompute para os dados do Flink. O exemplo a seguir mostra como criar uma delta table.
SET odps.sql.type.system.odps2=true; DROP TABLE mf_flink_upsert; CREATE TABLE mf_flink_upsert ( c1 int not null, c2 string, gt timestamp, primary key (c1) ) PARTITIONED BY (ds string) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; O conector Flink vem pré-carregado no Flink totalmente gerenciado, portanto, nenhuma instalação manual é necessária. Visualize os detalhes do conector no console do Realtime Compute for Apache Flink.
-
Crie uma tabela Flink, construa dados Flink em tempo real usando um job Flink SQL e implante o job após o desenvolvimento.
Na página de desenvolvimento de jobs Flink, crie e edite um job SQL. O exemplo abaixo define uma tabela de origem que gera dados aleatórios, uma tabela de resultados conectada ao MaxCompute e uma instrução INSERT para transferir os dados. Para mais informações sobre como desenvolver um job SQL, consulte Mapa de desenvolvimento de jobs.
-- Create a Flink source table. CREATE TEMPORARY TABLE fake_src_table ( c1 int, c2 VARCHAR, gt AS CURRENT_TIMESTAMP ) WITH ( 'connector' = 'faker', 'fields.c2.expression' = '#{superhero.name}', 'rows-per-second' = '100', 'fields.c1.expression' = '#{number.numberBetween ''0'',''1000''}' ); -- Create a temporary result table in Flink. CREATE TEMPORARY TABLE test_c_d_g ( c1 int, c2 VARCHAR, gt TIMESTAMP, ds varchar, PRIMARY KEY(c1) NOT ENFORCED ) PARTITIONED BY(ds) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_upsert', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj', 'upsert.write.bucket.num'='64' ); -- Flink computation logic. INSERT INTO test_c_d_g SELECT c1 AS c1, c2 AS c2, gt AS gt, date_format(gt, 'yyyyMMddHH') AS ds FROM fake_src_table;Parâmetros:
odps.end.point: Use o endpoint de rede interna da região correspondente.upsert.write.bucket.num: Este valor deve ser consistente com o valor da propriedade write.bucket.num da delta table criada no MaxCompute. -
Consulte a tabela MaxCompute para verificar se os dados foram gravados.
SELECT * FROM mf_flink_upsert WHERE ds=2023061517; -- Your results may differ because the source data is randomly generated. +------+----+------+----+ | c1 | c2 | gt | ds | +------+----+------+----+ | 0 | Skaar | 2023-06-16 01:59:41.116 | 2023061517 | | 21 | Supah Century | 2023-06-16 01:59:59.117 | 2023061517 | | 104 | Dark Gorilla Grodd | 2023-06-16 01:59:57.117 | 2023061517 | | 126 | Leader | 2023-06-16 01:59:39.116 | 2023061517 |
Apêndice: Parâmetros do conector Flink
-
Parâmetros básicos
Parâmetro
Obrigatório
Valor padrão
Descrição
connector
Sim
—
Defina o tipo de conector como
MaxCompute.odps.project.name
Sim
—
Nome do projeto MaxCompute.
odps.access.id
Sim
—
AccessKey ID da sua conta RAM. Consulte a página AccessKey Pair.
odps.access.key
Sim
—
AccessKey secret da sua conta RAM. Consulte a página AccessKey Pair.
odps.end.point
Sim
—
Endpoint do MaxCompute. Para obter uma lista de endpoints regionais, consulte Endpoints.
odps.tunnel.end.point
Não
—
Endpoint público do serviço Tunnel. Por padrão, as solicitações são roteadas automaticamente para o endpoint Tunnel apropriado. Defina este parâmetro para usar um endpoint específico e desativar o roteamento automático.
Para mais informações sobre os endpoints do Tunnel em diferentes regiões e redes, consulte Endpoints.
odps.tunnel.quota.name
Não
—
Nome da cota do Tunnel usada para acessar o MaxCompute.
table.name
Sim
—
Nome da tabela MaxCompute no formato
[project.][schema.]table.odps.namespace.schema
Não
false
Define se o modelo de três camadas deve ser usado. Para mais informações sobre o modelo de três camadas, consulte Operações de schema.
sink.operation
Sim
insert
Tipo de gravação. Os valores válidos são
insertouupsert.NotaO modo
upserté compatível apenas com delta tables do MaxCompute.sink.parallelism
Não
—
Paralelismo de gravação. Se não definido, assume o paralelismo da fonte upstream como padrão.
NotaGaranta que a propriedade da tabela
write.bucket.numseja um múltiplo inteiro do valor configurado para otimizar o desempenho de gravação e maximizar a economia de memória no nó Sink.sink.meta.cache.time
Não
400
Tamanho do cache de metadados.
sink.meta.cache.expire.time
Não
1200
Tempo de expiração do cache para metadados, em segundos.
sink.coordinator.enable
Não
true
Define se o modo coordenador deve ser ativado.
-
Parâmetros de partição
Parâmetro
Obrigatório
Valor padrão
Descrição
sink.partition
Não
—
Nome da partição onde os dados serão gravados.
Se você usar particionamento dinâmico, este parâmetro especifica o nome da partição pai das partições dinâmicas.
sink.partition.default-value
Não
__DEFAULT_PARTITION__
Nome padrão da partição quando o particionamento dinâmico é utilizado.
sink.dynamic-partition.limit
Não
100
Número máximo de partições que podem ser gravadas simultaneamente em um único checkpoint durante o particionamento dinâmico.
NotaAumentar significativamente este valor pode causar erros de falta de memória (OOM) no nó sink. Se o número de partições simultâneas exceder esse limite, o job falhará.
sink.group-partition.enable
Não
false
Define se o agrupamento por partição deve ser usado ao usar particionamento dinâmico.
sink.partition.assigner.class
Não
—
Classe de implementação
PartitionAssigner. -
Parâmetros de gravação no modo FileCached
Use o modo de cache de arquivo para jobs com um grande número de partições dinâmicas. Os parâmetros a seguir configuram esse modo.
Parâmetro
Obrigatório
Valor padrão
Descrição
sink.file-cached.enable
Não
false
Ativa o modo FileCached. Recomendado para jobs com muitas partições dinâmicas.
false: Modo FileCached desativado.
true: Modo FileCached ativado.
NotaQuando o número de partições dinâmicas for alto, use o modo de cache de arquivo.
sink.file-cached.tmp.dirs
Não
./local
Diretório padrão de cache de arquivos no modo FileCached.
sink.file-cached.writer.num
Não
16
Número de threads simultâneas de upload de dados para uma única tarefa no modo FileCached.
NotaNão aumente significativamente o valor deste parâmetro. Se um número excessivo de partições for gravado simultaneamente, erros OOM provavelmente ocorrerão.
sink.bucket.check-interval
Não
60000
Intervalo para verificação do tamanho do arquivo no modo FileCached. Unidade: milissegundos (ms).
sink.file-cached.rolling.max-size
Não
16 M
Tamanho máximo para um único arquivo de cache.
Quando um arquivo excede esse tamanho, ele é enviado.
sink.file-cached.memory
Não
64 M
Tamanho máximo de memória off-heap usado para gravação de arquivos no modo FileCached.
sink.file-cached.memory.segment-size
Não
128 KB
Tamanho do buffer usado para gravação de arquivos no modo FileCached.
sink.file-cached.flush.always
Não
true
Define se o cache deve ser usado para gravação de arquivos no modo FileCached.
sink.file-cached.write.max-retries
Não
3
Número de tentativas para upload de dados no modo FileCached.
-
Parâmetros de gravação
InsertouUpsertParâmetros de gravação Upsert
Parâmetro
Obrigatório
Valor padrão
Descrição
upsert.writer.max-retries
Não
3
Número de tentativas após falha de um writer upsert ao gravar dados em um bucket.
upsert.writer.buffer-size
Não
64 MB
Tamanho do cache para um único writer upsert no Flink.
NotaQuando a soma dos tamanhos de buffer de todos os buckets atinge o limiar predefinido, o sistema aciona automaticamente uma operação de flush para atualizar os dados no servidor.
Um writer upsert grava dados em vários buckets simultaneamente. Recomendamos aumentar o valor deste parâmetro para melhorar a eficiência da gravação.
Se os dados forem gravados em um grande número de partições, podem ocorrer erros OOM. Nesse caso, reduza o valor deste parâmetro.
upsert.writer.bucket.buffer-size
Não
1 MB
Tamanho do cache para um único bucket no Flink. Se os recursos de memória no servidor Flink forem insuficientes, reduza o valor deste parâmetro.
upsert.write.bucket.num
Sim
—
O número de buckets da tabela de destino deve ser igual ao valor de
write.bucket.num.upsert.write.slot-num
Não
1
Número de slots do Tunnel usados por uma única sessão.
upsert.commit.max-retries
Não
3
Número de tentativas para commit de uma sessão upsert.
upsert.commit.thread-num
Não
16
Paralelismo do commit de uma sessão upsert.
Não defina este valor muito alto. Um número elevado de commits simultâneos aumenta o consumo de recursos, o que pode causar problemas de desempenho ou uso excessivo de recursos.
upsert.major-compact.min-commits
Não
100
Número mínimo de commits necessários para acionar uma compactação principal (major compaction).
upsert.commit.timeout
Não
600
Tempo limite para o commit de uma sessão upsert. Unidade: segundos (s).
upsert.major-compact.enable
Não
false
Define se a compactação principal deve ser ativada.
upsert.flush.concurrent
Não
2
Número máximo de buckets nos quais os dados podem ser gravados simultaneamente em uma única partição.
NotaQuando os dados em um bucket passam por flush, um slot do Tunnel é ocupado.
NotaPara mais informações sobre configurações recomendadas de parâmetros para gravações upsert, consulte Configurações recomendadas de parâmetros para gravações upsert.
Parâmetros de gravação Insert
Parâmetro
Obrigatório
Valor padrão
Descrição
insert.commit.thread-num
Não
16
Paralelismo de uma sessão de commit.
insert.arrow-writer.enable
Não
false
Define se o formato Arrow deve ser usado.
insert.arrow-writer.batch-size
Não
512
Número máximo de linhas em um lote Arrow.
insert.arrow-writer.flush-interval
Não
100000
Intervalo de flush do writer. Unidade: milissegundos (ms).
insert.writer.buffer-size
Não
64 MB
Tamanho do cache do writer com buffer.