O conector do Hologres para Apache Flink permite gravar fluxos de dados de um cluster Flink open-source em tabelas do Hologres em tempo real. Esse conector oferece suporte tanto à interface Flink SQL quanto à DataStream API. Ele é open-source a partir do Apache Flink 1.11 e seus pacotes de release estão publicados no repositório Maven.
Pré-requisitos
Antes de começar, verifique se você possui:
Uma instância do Hologres com uma ferramenta de desenvolvimento conectada. Consulte Connect to HoloWeb
Um cluster Apache Flink (este tópico utiliza o Flink 1.15 em modo standalone). Para configurar um cluster, baixe o binário no site do Apache Flink e siga as instruções em Instalação local do Flink
-
Conectividade de rede entre o cluster Flink e sua instância do Hologres:
Mesma região: use o endpoint da Virtual Private Cloud (VPC) da sua instância do Hologres
Regiões diferentes: use o endpoint público
Adicione a dependência Maven
Inclua a dependência do conector do Hologres no seu arquivo pom.xml. Use a versão correspondente à sua instalação do Flink:
|
Versão do Apache Flink |
Conector do Hologres |
|
1.11 |
hologres-connector-flink-1.11:1.0.1 |
|
1.12 |
hologres-connector-flink-1.12:1.0.1 |
|
1.13 |
hologres-connector-flink-1.13:1.3.2 |
|
1.14 |
hologres-connector-flink-1.14:1.3.2 |
|
1.15 |
hologres-connector-flink-1.15:1.4.1 |
|
1.17 |
hologres-connector-flink-1.17:1.4.1 |
Use o Flink 1.15 ou superior para acessar mais recursos do conector.
O exemplo a seguir usa o Flink 1.15:
<dependency>
<groupId>com.alibaba.hologres</groupId>
<artifactId>hologres-connector-flink-1.15</artifactId>
<version>1.4.0</version>
<classifier>jar-with-dependencies</classifier>
</dependency>
Escolha um modo de gravação
O conector do Hologres oferece suporte a três modos de gravação. Escolha com base nos seus requisitos e no esquema da tabela:
|
Modo de gravação |
Opções |
Mais indicado para |
|
INSERT (padrão) |
|
Cargas de trabalho gerais em streaming |
|
Fixed copy |
|
Gravações com alto throughput |
|
Bulk load |
|
Gravações em estilo batch com máxima performance em tabelas sem chave primária ou em tabelas vazias com chave primária |
Observações sobre bulk load:
O bulk load reduz a carga da instância do Hologres em aproximadamente 66,7% em comparação ao fixed copy.
A gravação em uma tabela com chave primária causa bloqueios no nível da tabela. Para reduzir a granularidade do bloqueio para o nível de shard, defina
target-shards.enabled=true. Isso permite jobs de bulk load concorrentes.Se você usar bulk load para gravar em uma tabela com chave primária, a tabela deve estar vazia.
Requer Hologres V1.3.1+. O bulk load exige especificamente o Hologres V1.4.0+.
Comportamento de flush: Um lote de gravação é enviado (flushed) para o Hologres quando qualquer uma das condições abaixo for atendida:
O número de registros atinge
jdbcWriteBatchSize(padrão: 256)O tamanho dos dados em uma única thread atinge
jdbcWriteBatchByteSize(padrão: 2 MB)O tamanho total dos dados em todas as threads atinge
jdbcWriteBatchTotalByteSize(padrão: 20 MB)O tempo desde o último flush atinge
jdbcWriteFlushInterval(padrão: 10 segundos)
Grave dados usando Flink SQL
Use uma tabela sink do Flink SQL para gravar dados no Hologres. O tipo do conector é hologres.
CREATE TABLE sink (
user_id BIGINT,
user_name STRING,
price DECIMAL(38, 2),
sale_timestamp TIMESTAMP
) WITH (
'connector' = 'hologres',
'dbname' = '<your-database>',
'tablename' = '<your-table>',
'username' = '<your-access-key-id>',
'password' = '<your-access-key-secret>',
'endpoint' = '<your-endpoint>' -- Format: IP:Port
);
INSERT INTO sink SELECT * FROM source;
Substitua os placeholders pelos valores reais:
|
Placeholder |
Descrição |
|
|
Nome do banco de dados do Hologres |
|
|
Nome da tabela de destino no Hologres |
|
|
Seu AccessKey ID da Alibaba Cloud. Obtenha-o na página AccessKey Pair |
|
|
Seu AccessKey secret da Alibaba Cloud. Obtenha-o na página AccessKey Pair |
|
|
Endpoint VPC da sua instância do Hologres, no formato |
Grave dados usando a DataStream API
As aplicações de exemplo a seguir demonstram padrões comuns. Todos os exemplos gravam no Hologres usando as opções do conector descritas em Referência de opções do conector.
|
Aplicação |
Descrição |
|
|
Grava dados usando a interface Flink SQL |
|
|
Converte um DataStream em uma tabela e depois grava usando a interface Flink SQL |
|
|
Grava fluxos de dados diretamente usando a interface Flink DataStream |
|
|
Conta visitantes únicos (UVs) em tempo real usando roaring bitmaps e tabelas de dimensão do Hologres, gravando os resultados no Hologres |
|
|
Particiona dados por shard antes da gravação, usando a interface DataStream. Adequado para importação em massa de dados em múltiplas tabelas vazias com chaves primárias; produz um efeito semelhante ao INSERT OVERWRITE |
Referência de opções do conector
Opções obrigatórias
|
Opção |
Descrição |
|
|
Tipo do conector sink. Defina como |
|
|
Nome do banco de dados do Hologres |
|
|
Nome da tabela de destino no Hologres |
|
|
Seu AccessKey ID |
|
|
Seu AccessKey secret |
|
|
Endpoint VPC da sua instância do Hologres, no formato |
Opções de conexão
|
Opção |
Padrão |
Descrição |
|
|
3 |
Número de conexões Java Database Connectivity (JDBC) por pool de conexões por tarefa do Flink. Aumente proporcionalmente aos requisitos de throughput |
|
|
(nenhum) |
Nome do pool de conexões. Tabelas que compartilham um pool devem ter o mesmo valor de |
|
|
false |
Quando definido como |
|
|
10 |
Número máximo de tentativas de reconexão em caso de falha |
|
|
1000 |
Intervalo base de nova tentativa em milissegundos. Intervalo de retry = |
|
|
5000 |
Incremento gradual para intervalos de retry em milissegundos |
|
|
60000 |
Tempo limite de ociosidade para conexões de gravação e point query em milissegundos. A conexão é liberada após o término do tempo limite |
|
|
60000 |
Time-to-live (TTL) do cache de esquema da tabela em milissegundos |
|
|
-1 |
Fator que aciona a atualização automática do cache. O cache é atualizado quando seu TTL restante cai abaixo de |
|
|
disable |
Modo de transmissão criptografada por SSL. Valores válidos: |
|
|
(nenhum) |
Caminho para o arquivo de certificado CA no cluster Flink. Obrigatório quando |
|
|
false |
Quando definido como |
Opções de sink
|
Opção |
Padrão |
Descrição |
|
|
insertorignore |
Modo de gravação para o sink. Consulte Streaming semantics |
|
|
true |
Quando definido como |
|
|
false |
Se definido como |
|
|
false |
Quando definido como |
|
|
256 |
Número máximo de registros por lote de gravação por nó de sink em streaming |
|
|
2097152 |
Tamanho máximo de dados por lote de gravação por thread, em bytes (padrão: 2 MB) |
|
|
20971520 |
Tamanho máximo de dados por lote de gravação em todas as threads, em bytes (padrão: 20 MB) |
|
|
10000 |
Tempo máximo de espera antes de enviar (flush) um lote de gravação, em milissegundos (padrão: 10 segundos) |
|
|
false |
Controla a sintaxe SQL usada para gravações. |
|
|
true |
Quando definido como |
|
|
false |
Quando definido como |
|
|
true |
Se definido como |
|
|
false |
Quando definido como |
|
|
false |
Quando definido como |
|
|
binary |
Formato do protocolo quando |
|
|
false |
Quando definido como |
|
|
false |
Quando definido como |
Opções de point query
Estas opções aplicam-se ao usar o Hologres como uma tabela de dimensão do Flink para lookup joins.
|
Opção |
Padrão |
Descrição |
|
|
128 |
Número máximo de solicitações por lote por thread para point queries em tabelas de dimensão |
|
|
256 |
Número máximo de solicitações enfileiradas por thread |
|
|
false |
Quando definido como |
|
|
None |
Política de cache para lookups de tabela de dimensão. |
|
|
10000 |
Número máximo de registros em cache. Aplica-se apenas quando |
|
|
(sem expiração) |
Intervalo de atualização do cache em milissegundos. Aplica-se apenas quando |
|
|
true |
Quando definido como |
Mapeamentos de tipos de dados
Para mapeamentos de tipos de dados entre Apache Flink e Hologres, consulte a seção "Mapeamentos de tipos de dados entre Realtime Compute for Apache Flink ou Blink e Hologres" em Data types.
Próximos passos
Create a Hologres result table — aprenda sobre semântica de streaming e a opção
mutatetypeData types — referência completa de mapeamento de tipos de dados