Importe dados em lote para o Hologres pelo Flink connector para uma ingestão de dados altamente eficiente e de baixa carga.
Contexto
O Hologres integra-se ao Apache Flink para oferecer recursos de streaming de dados em tempo real. Em casos de uso que não são sensíveis a tempo — como carregar dados históricos, processar dados offline ou agregar logs — a importação em lote é a abordagem recomendada. Esse método grava grandes volumes de dados no Hologres de uma só vez, o que é mais eficiente e economiza recursos. Escolha entre importação em tempo real e importação em lote conforme as necessidades do seu negócio e os recursos disponíveis. Para mais informações sobre importação em tempo real, consulte Realtime Compute for Apache Flink.
Pré-requisitos
Você adquiriu uma instância do Hologres. Para mais informações, consulte Purchase a Hologres instance.
-
Você implantou um cluster Apache Flink, versão 1.15 ou posterior. Para mais informações, consulte os tópicos a seguir:
Apache Flink: Deploy Flink.
Realtime Compute for Apache Flink: Activate Realtime Compute for Apache Flink.
Importação em lote com o Realtime Compute for Apache Flink
-
Crie uma tabela de resultado no Hologres para armazenar os dados importados do Flink. Para as instruções, consulte Connect to HoloWeb and run queries. Este tópico usa a tabela
test_sink_customercomo exemplo.-- Create a Hologres result table. CREATE TABLE test_sink_customer ( c_custkey BIGINT, c_name TEXT, c_address TEXT, c_nationkey INT, c_phone TEXT, c_acctbal NUMERIC(15,2), c_mktsegment TEXT, c_comment TEXT, "date" DATE ) WITH ( distribution_key="c_custkey,date", orientation="column" );NotaOs nomes de campo e os tipos de dados na tabela de origem do Flink devem corresponder aos da tabela de resultado do Hologres.
-
Faça login no Realtime Compute for Apache Flink console. Na página Deployments, clique em Create Deployment. Configure os parâmetros do deployment e clique em Deploy. Para mais informações sobre os parâmetros, consulte Deploy a JAR job.
A tabela a seguir descreve os principais parâmetros.
Parâmetro
Descrição
Deployment Type
Selecione JAR.
Deployment Mode
Escolha entre modo stream ou modo batch. Este tópico usa o modo batch como exemplo.
Engine Version
Para mais informações sobre as versões da engine, consulte Engine versions e Lifecycle policies. Este tópico usa a versão
vvr-8.0.7-flink-1.17como exemplo.JAR URI
Faça upload do Flink connector open source: hologres-connector-flink-repartition.jar.
NotaUse o Flink connector open source para importar dados em lote para o Hologres. Para o código-fonte do Flink connector, consulte o repositório oficial do Hologres no GitHub.
Entry Point Class
A classe de entrada do programa. A classe principal do Flink connector é
com.alibaba.ververica.connectors.hologres.example.FlinkToHoloRePartitionExample.Entry Point Main Arguments
Informe o caminho para o arquivo
repartition.sql. No runtime do Realtime Compute for Apache Flink, os arquivos de dependência adicionais ficam armazenados em/flink/usrlib/. Portanto, o argumento completo é--sqlFilePath="/flink/usrlib/repartition.sql".Additional Dependencies
Faça upload do arquivo
repartition.sql. Trata-se de um script Flink SQL usado para definir a fonte de dados, declarar a tabela de resultado e configure a conexão com o Hologres. O código a seguir é um exemplo do arquivorepartition.sql.-- DDL for the source table. This example uses the Flink DataGen connector to generate test data. CREATE TEMPORARY TABLE source_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ) WITH ( 'connector' = 'datagen' ,'rows-per-second' = '10000' ,'number-of-rows' = '1000000' ); -- DQL for the source table. The query result must match the schema of the result table defined in the sink DDL, including the number and types of fields. SELECT *, cast('2024-04-21' as DATE) FROM source_table; -- DDL for the sink table. This declares the result table and configures the connection to Hologres. CREATE TABLE sink_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ,`date` DATE ) WITH ( 'connector' = 'hologres' ,'dbname' = 'doc_****' ,'tablename' = 'test_sink_customer' ,'username' = 'yourAccessKeyId' ,'password' = 'yourAccessKeySecret' ,'endpoint' = 'hgpostcn-cn-7pp2e1k7****-cn-hangzhou.hologres.aliyuncs.com:80' ,'jdbccopywritemode' = 'true' ,'bulkload' = 'true' ,'target-shards.enabled'='true' );NotaPara mais informações sobre os parâmetros de conexão do Hologres no arquivo
repartition.sql, consulte Hologres Flink connector parameters. -
Clique em no nome do deployment e acesse a página Deployment Details. Na seção Resource Configurations, modifique o Parallelism.
NotaRecomendamos definir o paralelismo como o ShardCount da tabela de resultado do Hologres.
-
Consulte a tabela de resultado do Hologres.
Após o envio do job do Flink, consulte os dados gravados no Hologres. Instrução de exemplo:
SELECT * FROM test_sink_customer;
Importação em lote com o Apache Flink
-
Crie uma tabela de resultado no Hologres para armazenar os dados importados do Flink. Para as instruções, consulte Connect to HoloWeb and run queries. Este tópico usa a tabela
test_sink_customercomo exemplo.-- Create a Hologres result table. CREATE TABLE test_sink_customer ( c_custkey BIGINT, c_name TEXT, c_address TEXT, c_nationkey INT, c_phone TEXT, c_acctbal NUMERIC(15,2), c_mktsegment TEXT, c_comment TEXT, "date" DATE ) WITH ( distribution_key="c_custkey,date", orientation="column" );NotaDefina a quantidade de shards conforme o volume dos seus dados. Para mais informações sobre shards, consulte Manage table groups and shard count.
-
Crie o arquivo
repartition.sqle faça upload dele para qualquer diretório do ambiente do seu cluster Flink. Este tópico usa o caminho/flink-1.15.4/src/repartition.sqlcomo exemplo. O código a seguir é um exemplo do arquivorepartition.sql.NotaTrata-se de um script Flink SQL usado para definir a fonte de dados, declarar a tabela de resultado e configure a conexão com o Hologres.
-- DDL for the source table. This example uses the Flink DataGen connector to generate test data. CREATE TEMPORARY TABLE source_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ) WITH ( 'connector' = 'datagen' ,'rows-per-second' = '10000' ,'number-of-rows' = '1000000' ); -- DQL for the source table. The query result must match the schema of the result table defined in the sink DDL, including the number and types of fields. SELECT *, cast('2024-04-21' as DATE) FROM source_table; -- DDL for the sink table. This declares the result table and configures the connection to Hologres. CREATE TABLE sink_table ( c_custkey BIGINT ,c_name STRING ,c_address STRING ,c_nationkey INTEGER ,c_phone STRING ,c_acctbal NUMERIC(15, 2) ,c_mktsegment STRING ,c_comment STRING ,`date` DATE ) WITH ( 'connector' = 'hologres' ,'dbname' = 'doc_****' ,'tablename' = 'test_sink_customer' ,'username' = 'yourAccessKeyId' ,'password' = 'yourAccessKeySecret' ,'endpoint' = 'hgpostcn-cn-7pp2e1k7****-cn-hangzhou.hologres.aliyuncs.com:80' ,'jdbccopywritemode' = 'true' ,'bulkload' = 'true' ,'target-shards.enabled'='true' );A tabela a seguir descreve os principais parâmetros.
Parâmetro
Obrigatório
Descrição
connector
Sim
O tipo do connector. O valor deve ser
hologres.dbname
Sim
O nome do banco de dados do Hologres.
tablename
Sim
O nome da tabela do Hologres que receberá os dados.
username
Sim
O AccessKey ID da sua conta Alibaba Cloud.
Obtenha o AccessKey ID na página AccessKey Pair.
password
Sim
O AccessKey secret correspondente ao seu AccessKey ID.
endpoint
Sim
O endpoint VPC da instância do Hologres. Acesse a página de detalhes da instância no Hologres console e obtenha o endpoint na seção Configurations.
NotaO endpoint deve incluir o número da porta no formato
ip:port. Use o endpoint VPC para conexões dentro da mesma região. Use o endpoint público para conexões entre regiões.jdbccopywritemode
Não
O método de gravação de dados. Valores válidos:
-
false(padrão): usa o métodoINSERT. -
true: usa o métodoCOPY, que inclui streamingCOPY(Fixed Copy) e batchCOPY. Por padrão, o streamingCOPYé utilizado.NotaEm comparação com o método
INSERT, o streamingCOPYadota um modelo de streaming para alcançar throughput maior, menor latência de dados e menor consumo de memória no cliente, já que os dados não são acumulados em lotes. No entanto, não há suporte para retração de dados.
bulkload
Não
Determina se o método batch
COPYserá usado. Valores válidos:-
true: usa batchCOPY. Essa configuração só tem efeito quandojdbccopywritemodetambém está definido comotrue. Caso contrário, o streamingCOPYé utilizado.Nota-
Em comparação com o streaming
COPY, o batchCOPYé mais eficiente e utiliza os recursos do Hologres de forma mais eficaz, resultando em desempenho de gravação superior. Escolha o método de gravação apropriado conforme os requisitos do seu negócio. -
Ao usar batch
COPYpara gravar em uma tabela com chave primária, podem ocorrer bloqueios de tabela. Defina o parâmetrotarget-shards.enabledcomotruepara reduzir a granularidade do bloqueio do nível de tabela para o nível de shard. Assim, várias tarefas de importação em lote podem executar simultaneamente e a contenção de bloqueio de tabela é reduzida. Em relação ao streamingCOPY, essa abordagem reduz significativamente a carga sobre a instância do Hologres ao gravar em uma tabela com chave primária. Testes mostram uma redução de carga de aproximadamente 66,7%. -
Ao usar batch
COPY, se a tabela de destino tiver chave primária, ela deverá estar vazia antes da operação de gravação. Caso contrário, o processo de gravação será desacelerado pela deduplicação de dados baseada na chave primária.
-
-
false(padrão): não usa batchCOPY.
target-shards.enabled
Não
Determina se a gravação em lote por target shard será ativada. Valores válidos:
-
true: ative a gravação em lote por target shard. Quando os dados de origem são reparticionados por shard, isso reduz a granularidade do bloqueio para o nível de shard. -
false(padrão): desativa este recurso.
NotaPara mais informações sobre os parâmetros de conexão do Hologres no arquivo
repartition.sql, consulte Hologres Flink connector parameters. -
-
No ambiente do seu cluster Flink, faça upload do Flink connector open source hologres-connector-flink-repartition.jar para qualquer diretório. Este tópico usa o diretório raiz como exemplo.
NotaUse o Flink connector open source para importar dados em lote para o Hologres. Para o código-fonte do Flink connector, consulte o repositório oficial do Hologres no GitHub.
-
Envie o job do Flink. Comando de exemplo:
./bin/flink run -Dexecution.runtime-mode=BATCH -p 3 -c com.alibaba.ververica.connectors.hologres.example.FlinkToHoloRePartitionExample hologres-connector-flink-repartition.jar --sqlFilePath="/flink-1.15.4/src/repartition.sql"Parâmetros do comando acima:
Dexecution.runtime-mode: o modo de execução do job do Flink. Para mais informações, consulte Execution Mode.p: o paralelismo do job. Recomendamos definir esse valor como o ShardCount da tabela de resultado, ou como um divisor do ShardCount.c: o nome totalmente qualificado da classe principal no arquivo hologres-connector-flink-repartition.jar.sqlFilePath: o caminho para o arquivorepartition.sql.
-
Consulte a tabela de resultado do Hologres.
Após o envio do job do Flink, consulte os dados gravados no Hologres. Instrução de exemplo:
SELECT * FROM test_sink_customer;