Todos os produtos
Search
Central de documentação

Hologres:Streaming writes to Hologres with Apache Flink 1.11+

Última atualização: Sep 18, 2026

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)

jdbcCopyWriteMode=false

Cargas de trabalho gerais em streaming

Fixed copy

jdbcCopyWriteMode=true, bulkLoad=false

Gravações com alto throughput

Bulk load

jdbcCopyWriteMode=true, bulkLoad=true

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

<your-database>

Nome do banco de dados do Hologres

<your-table>

Nome da tabela de destino no Hologres

<your-access-key-id>

Seu AccessKey ID da Alibaba Cloud. Obtenha-o na página AccessKey Pair

<your-access-key-secret>

Seu AccessKey secret da Alibaba Cloud. Obtenha-o na página AccessKey Pair

<your-endpoint>

Endpoint VPC da sua instância do Hologres, no formato IP:Port. Encontre-o na página de detalhes da instância no console do Hologres

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

FlinkSQLToHoloExample

Grava dados usando a interface Flink SQL

FlinkDSAndSQLToHoloExample

Converte um DataStream em uma tabela e depois grava usando a interface Flink SQL

FlinkDataStreamToHoloExample

Grava fluxos de dados diretamente usando a interface Flink DataStream

FlinkRoaringBitmapAggJob

Conta visitantes únicos (UVs) em tempo real usando roaring bitmaps e tabelas de dimensão do Hologres, gravando os resultados no Hologres

FlinkToHoloRePartitionExample

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

connector

Tipo do conector sink. Defina como hologres

dbname

Nome do banco de dados do Hologres

tablename

Nome da tabela de destino no Hologres

username

Seu AccessKey ID

password

Seu AccessKey secret

endpoint

Endpoint VPC da sua instância do Hologres, no formato IP:Port. Use o endpoint VPC para acesso na mesma região; use o endpoint público para acesso entre regiões

Opções de conexão

Opção

Padrão

Descrição

connectionSize

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

connectionPoolName

(nenhum)

Nome do pool de conexões. Tabelas que compartilham um pool devem ter o mesmo valor de connectionSize. Por padrão, cada tabela usa seu próprio pool

fixedConnectionMode

false

Quando definido como true, operações de gravação e point query não consomem conexões. Requer conector 1.2.0+ e Hologres 1.3+. Este recurso está em beta

jdbcRetryCount

10

Número máximo de tentativas de reconexão em caso de falha

jdbcRetrySleepInitMs

1000

Intervalo base de nova tentativa em milissegundos. Intervalo de retry = jdbcRetrySleepInitMs + (contagem de retries x jdbcRetrySleepStepMs)

jdbcRetrySleepStepMs

5000

Incremento gradual para intervalos de retry em milissegundos

jdbcConnectionMaxIdleMs

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

jdbcMetaCacheTTL

60000

Time-to-live (TTL) do cache de esquema da tabela em milissegundos

jdbcMetaAutoRefreshFactor

-1

Fator que aciona a atualização automática do cache. O cache é atualizado quando seu TTL restante cai abaixo de jdbcMetaCacheTTL / jdbcMetaAutoRefreshFactor. O valor -1 desativa a atualização automática

connection.ssl.mode

disable

Modo de transmissão criptografada por SSL. Valores válidos: disable, require, verify-ca, verify-full

connection.ssl.root-cert.location

(nenhum)

Caminho para o arquivo de certificado CA no cluster Flink. Obrigatório quando connection.ssl.mode é verify-ca ou verify-full

jdbcDirectConnect

false

Quando definido como true, o conector verifica se o Flink pode se conectar diretamente aos nós Frontend (FE) do Hologres no ambiente atual e usa conexões diretas se disponíveis

Opções de sink

Opção

Padrão

Descrição

mutatetype

insertorignore

Modo de gravação para o sink. Consulte Streaming semantics

ignoredelete

true

Quando definido como true, mensagens de retração (solicitações DELETE geradas por operações GROUP BY do Flink) são ignoradas. Aplica-se apenas quando a semântica de streaming é utilizada

createparttable

false

Se definido como true, partições são criadas automaticamente com base nos valores da chave de partição ao gravar em uma tabela particionada. Use com cautela: valores inválidos de chave de partição criam partições inválidas

ignoreNullWhenUpdate

false

Quando definido como true, valores nulos não são gravados durante uma atualização se mutatetype='insertOrUpdate'

jdbcWriteBatchSize

256

Número máximo de registros por lote de gravação por nó de sink em streaming

jdbcWriteBatchByteSize

2097152

Tamanho máximo de dados por lote de gravação por thread, em bytes (padrão: 2 MB)

jdbcWriteBatchTotalByteSize

20971520

Tamanho máximo de dados por lote de gravação em todas as threads, em bytes (padrão: 20 MB)

jdbcWriteFlushInterval

10000

Tempo máximo de espera antes de enviar (flush) um lote de gravação, em milissegundos (padrão: 10 segundos)

jdbcUseLegacyPutHandler

false

Controla a sintaxe SQL usada para gravações. true: INSERT INTO xxx(c0,c1,...) VALUES (?,...) ON CONFLICT. false: INSERT INTO xxx(c0,c1,...) SELECT unnest(?), unnest(?),... ON CONFLICT

jdbcEnableDefaultForNotNullColumn

true

Quando definido como true, valores nulos gravados em colunas NOT NULL sem valor padrão são substituídos por: "" (STRING), 0 (NUMBER) ou 1970-01-01 00:00:00 (DATE/TIMESTAMP/TIMESTAMPTZ)

remove-u0000-in-text.enabled

false

Quando definido como true, substitui caracteres u0000 em colunas TEXT que não estejam codificadas em UTF-8

deduplication.enabled

true

Se definido como true, apenas o último registro com cada valor de chave primária é retido quando múltiplos registros compartilham a mesma chave no mesmo lote de gravação. Quando false, todos os registros são gravados sem deduplicação. Desativar a deduplicação pode causar falhas de gravação se todos os registros em um lote compartilharem a mesma chave primária

aggressive.enabled

false

Quando definido como true, o conector confirma os dados imediatamente ao detectar uma conexão ociosa, independentemente do tamanho do lote. Reduz a latência de gravação sob baixo tráfego

jdbcCopyWriteMode

false

Quando definido como true, usa o protocolo COPY em vez de INSERT. Consulte Escolha um modo de gravação. Requer Hologres V1.3.1+

jdbcCopyWriteFormat

binary

Formato do protocolo quando jdbcCopyWriteMode=true. O formato binary é mais rápido; text está disponível como alternativa. Requer Hologres V1.3.1+

bulkLoad

false

Quando definido como true e jdbcCopyWriteMode=true, usa bulk load em vez de fixed copy. Requer Hologres V1.4.0+

target-shards.enabled

false

Quando definido como true, o bulk load visa shards individuais, reduzindo a granularidade do bloqueio e permitindo jobs de bulk load concorrentes. Requer Hologres V1.4.1+

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

jdbcReadBatchSize

128

Número máximo de solicitações por lote por thread para point queries em tabelas de dimensão

jdbcReadBatchQueueSize

256

Número máximo de solicitações enfileiradas por thread

async

false

Quando definido como true, o conector processa múltiplas solicitações e respostas simultaneamente, melhorando o throughput da consulta. Não há garantia de que as solicitações sejam concluídas em ordem

cache

None

Política de cache para lookups de tabela de dimensão. None: sem cache. LRU: armazena em cache um subconjunto dos dados da tabela de dimensão na memória

cachesize

10000

Número máximo de registros em cache. Aplica-se apenas quando cache=LRU

cachettlms

(sem expiração)

Intervalo de atualização do cache em milissegundos. Aplica-se apenas quando cache=LRU. Por padrão, as entradas em cache não expiram

cacheempty

true

Quando definido como true, os resultados de consultas JOIN que não retornam linhas também são armazenados em cache

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