Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector do ClickHouse

Última atualização: Aug 20, 2026

Este tópico descreve o uso do conector do ClickHouse para gravar dados no ClickHouse.

Informações básicas

O ClickHouse é um sistema de gerenciamento de banco de dados orientado a colunas para Processamento Analítico Online (OLAP). Para mais informações, consulte O que é o ClickHouse?.

A tabela a seguir descreve as capacidades do conector do ClickHouse.

Categoria

Detalhes

Tipo compatível

Apenas tabela de resultados

Modo de execução

Modos em lote e streaming

Formato de dados

Não aplicável

Métricas específicas do conector

  • numRecordsOut

  • numRecordsOutPerSecond

  • currentSendTime

Nota

Para mais informações sobre essas métricas, consulte Monitoring metrics.

Tipo de API

SQL

Suporte para atualizar ou excluir dados em uma tabela de resultados

O conector permite atualizações e exclusões se você especificar uma chave primária na DDL da tabela de resultados do Flink e definir o parâmetro ignoreDelete como false. No entanto, essa opção reduz significativamente o desempenho.

Recursos

  • Grave dados diretamente nas tabelas locais de uma tabela distribuída do ClickHouse.

  • Garanta semântica exactly-once ao gravar no ClickHouse no Alibaba Cloud E-MapReduce (EMR).

Pré-requisitos

  • Crie uma tabela no ClickHouse. Para mais informações, consulte Criar uma nova tabela.

  • Configure uma lista de permissões.

    • Se você utiliza o ApsaraDB for ClickHouse, consulte Configure a whitelist.

    • Para usuários do ClickHouse no Alibaba Cloud E-MapReduce (EMR), consulte Manage security groups.

    • Caso utilize um cluster ClickHouse autogerenciado em uma instância ECS, consulte Security group overview.

    • Em outras configurações, configure a lista de permissões na máquina onde o ClickHouse está implantado para permitir o acesso a partir da implantação do Realtime Compute for Apache Flink.

    Nota

    Para visualizar o segmento de rede do vSwitch do Realtime Compute for Apache Flink, consulte How do I configure a whitelist?.

Limitações

  • O parâmetro sink.parallelism não é compatível.

  • Por padrão, o sink do ClickHouse oferece semântica at-least-once.

  • O conector do ClickHouse é compatível apenas com o Ververica Runtime (VVR) 3.0.2 e versões posteriores.

  • A opção ignoreDelete está disponível somente nas versões VVR 3.0.3, VVR 4.0.7 e posteriores.

  • O tipo de dados Nested do ClickHouse é compatível apenas a partir do VVR 4.0.10.

  • Gravações diretas nas tabelas locais de uma tabela distribuída são permitidas exclusivamente no VVR 4.0.11 e versões mais recentes.

  • A semântica exactly-once para gravação no ClickHouse no EMR está disponível apenas no VVR 4.0.11 e posteriores. Contudo, devido a alterações nas capacidades do product, essa semântica não está disponível para o ClickHouse no EMR v3.45.1 e em versões posteriores à v5.11.1.

  • O modo de gravação balance, que distribui dados uniformemente entre os nós das tabelas locais, está disponível apenas no VVR 8.0.7 e versões posteriores.

  • A gravação em uma tabela local do ClickHouse é compatível somente com a Edição Compatível com a Comunidade do ApsaraDB for ClickHouse.

  • Ao registrar um catálogo do ClickHouse, o nome do banco de dados padrão não pode conter hífen (-). Caso contrário, a validação da URL JDBC falhará.

Sintaxe

CREATE TABLE clickhouse_sink (
  id INT,
  name VARCHAR,
  age BIGINT,
  rate FLOAT
) WITH (
  'connector' = 'clickhouse',
  'url' = '<yourUrl>',
  'userName' = '<yourUsername>',
  'password' = '<yourPassword>',
  'tableName' = '<yourTablename>',
  'maxRetryTimes' = '3',
  'batchSize' = '8000',
  'flushIntervalMs' = '1000',
  'ignoreDelete' = 'true',
  'shardWrite' = 'false',
  'writeMode' = 'partition',
  'shardingKey' = 'id'
);

Parâmetros WITH

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

connector

Tipo da tabela de resultados.

String

Sim

N/A

Defina o valor como clickhouse.

url

URL JDBC do seu cluster ClickHouse.

String

Sim

N/A

O formato da URL é jdbc:clickhouse://<yourNetworkAddress>:<port>/<yourDatabaseName>. Ao gravar diretamente em uma tabela local, obtenha o endereço IP do nó executando select * from system.clusters no ClickHouse. Se nenhum nome de banco de dados for especificado, o sistema usará o banco de dados default.

Nota

Para gravar dados em uma tabela distribuída do ClickHouse, a url deve ser a URL JDBC de um nó no cluster que hospeda a tabela distribuída.

userName

Nome de usuário para acessar o ClickHouse.

String

Sim

N/A

N/A

password

Senha para acessar o ClickHouse.

String

Sim

N/A

N/A

tableName

Nome da tabela do ClickHouse.

String

Sim

N/A

N/A

maxRetryTimes

Número máximo de tentativas após falha na inserção de dados na tabela de resultados.

Int

Não

3

N/A

batchSize

Quantidade de registros a serem gravados em um único lote.

Int

Não

100

Quando a quantidade de entradas de dados no cache atingir o valor do parâmetro batchSize ou o tempo de espera exceder flushIntervalMs, o sistema grava automaticamente os dados do cache na tabela do ClickHouse.

flushIntervalMs

Intervalo de tempo, em milissegundos, para liberar o buffer.

Long

Não

1000

A unidade é milissegundos.

ignoreDelete

Define se as mensagens de exclusão devem ser ignoradas.

Boolean

Não

true

Valores válidos:

  • true (padrão): Ignora mensagens de exclusão.

  • false: Não ignora mensagens de exclusão.

    Se esta opção estiver definida como false e uma chave primária for definida na DDL, o conector usará uma instrução ALTER para excluir dados no ClickHouse.

Nota

Se você definir ignoreDelete como false, não poderá usar o modo de gravação partition.

shardWrite

Para uma tabela distribuída do ClickHouse, define se os dados devem ser gravados diretamente nas tabelas locais subjacentes.

Boolean

Não

false

Valores válidos:

  • false (padrão): O conector grava dados na tabela distribuída, que então roteia os dados para as tabelas locais correspondentes. O parâmetro tableName deve ser o nome da tabela distribuída.

  • true: O conector grava dados diretamente nas tabelas locais, ignorando a tabela distribuída.

    Esta configuração é recomendada para melhorar o throughput de gravação.

    • Para especificar manualmente em quais tabelas locais gravar na url, o tableName deve ser o nome da tabela local. Exemplo:

      'url' = 'jdbc:clickhouse://192.XX.XX.1:3002,192.XX.XX.2:3002/default'
          'tableName' = 'local_table'
    • Caso prefira não especificar os nós manualmente, defina o parâmetro inferLocalTable como true para permitir que o Flink descubra automaticamente os nós da tabela local. Nesse cenário, tableName deve ser o nome da tabela distribuída e url deve ser a URL JDBC de um nó no cluster. Exemplo:

      'url' = 'jdbc:clickhouse://192.XX.XX.1:3002/default' // The JDBC URL of a node in the cluster.
          'tableName' = 'distribute_table'

inferLocalTable

Ao gravar em uma tabela distribuída do ClickHouse, define se o sistema deve descobrir automaticamente as tabelas locais subjacentes e gravar diretamente nelas.

Boolean

Não

false

Valores válidos:

  • false (padrão): Se você estiver gravando em uma tabela distribuída e especificar apenas um nó no parâmetro url, o conector não tentará descobrir as tabelas locais. Ele gravará na tabela distribuída, que então roteará os dados para as tabelas locais.

  • true: O Flink tenta descobrir as tabelas locais e grava diretamente nelas. Isso requer que shardWrite esteja definido como true, que tableName seja o nome da tabela distribuída e que url seja a URL JDBC de um nó no cluster.

Nota

Este parâmetro é ignorado ao gravar em uma tabela não distribuída.

writeMode

Estratégia para gravar dados nas tabelas locais de uma tabela distribuída.

Enum

Não

default

Valores válidos:

  • default (padrão): Sempre grava na tabela local do primeiro nó especificado na URL.

  • partition: Distribui dados com base em uma chave de fragmentação, garantindo que registros com a mesma chave sejam gravados no mesmo nó da tabela local.

  • random: Grava dados em um nó da tabela local selecionado aleatoriamente.

  • balance: Utiliza uma estratégia round-robin para distribuir dados uniformemente por todos os nós da tabela local.

Nota

Se você definir writeMode como partition, deverá definir ignoreDelete como true.

shardingKey

Chave usada para particionar dados entre os nós da tabela local.

String

Não

N/A

Quando writeMode estiver definido como 'partition', o parâmetro shardingKey é obrigatório. Pode conter vários campos separados por vírgulas (,).

exactlyOnce

Define se a semântica exactly-once deve ser ativada.

Boolean

Não

false

Valores válidos:

  • true: Ativa a semântica exactly-once.

  • false (padrão): Desativa a semântica exactly-once.

Nota
  • A semântica exactly-once é compatível apenas com gravação no ClickHouse no EMR. Defina este parâmetro como true somente se o seu sink for um cluster ClickHouse no EMR.

  • A semântica exactly-once não é compatível com gravação em tabelas locais do ClickHouse com a estratégia de partição. Portanto, se exactlyOnce estiver definido como true, writeMode não poderá ser definido como partition.

Mapeamento de tipos de dados

Tipo Flink

Tipo ClickHouse

BOOLEAN

UInt8 / Boolean

Nota

O ClickHouse v21.12 e versões posteriores suportam o tipo Boolean. Em versões anteriores, o tipo BOOLEAN do Flink é mapeado para o tipo UInt8 do ClickHouse.

TINYINT

Int8

SMALLINT

Int16

INTEGER

Int32

BIGINT

Int64

BIGINT

UInt32

FLOAT

Float32

DOUBLE

Float64

CHAR

FixedString

VARCHAR

String

BINARY

FixedString

VARBINARY

String

DATE

Date

TIMESTAMP(0)

DateTime

TIMESTAMP(x)

Datetime64(x)

DECIMAL

DECIMAL

ARRAY

ARRAY

Nested

Nota

O conector do ClickHouse não suporta os tipos TIME, MAP, MULTISET e ROW do Flink.

Para usar o tipo de dados Nested do ClickHouse, mapeie-o para um tipo ARRAY no Flink. Por exemplo:

-- ClickHouse
CREATE TABLE visits (
  StartDate Date,
  Goals Nested
  (
    ID UInt32,
    OrderID String
  )
  ...
);

Faça o mapeamento dos tipos da seguinte forma:

-- Flink
CREATE TABLE visits (
  StartDate DATE,
  `Goals.ID` ARRAY<LONG>,
  `Goals.OrderID` ARRAY<STRING>
);
Nota

Em versões do VVR anteriores à 6.0.6, a gravação de dados DateTime64 com o driver JDBC oficial do ClickHouse causava perda de precisão ao truncar valores para segundos. Como resultado, apenas dados TIMESTAMP com precisão de segundos, como TIMESTAMP(0), podiam ser gravados com sucesso. O VVR 6.0.6 e versões posteriores resolvem esse problema, permitindo gravações precisas de dados DateTime64.

Exemplos

  • Exemplo 1: Gravar em uma tabela de nó único.

    CREATE TEMPORARY TABLE clickhouse_source (
      id INT,
      name VARCHAR,
      age BIGINT,
      rate FLOAT
    ) WITH (
      'connector' = 'datagen',
      'rows-per-second' = '50'
    );
    
    CREATE TEMPORARY TABLE clickhouse_output (
      id INT,
      name VARCHAR,
      age BIGINT,
      rate FLOAT
    ) WITH (
      'connector' = 'clickhouse',
      'url' = '<yourUrl>',
      'userName' = '<yourUsername>',
      'password' = '<yourPassword>',
      'tableName' = '<yourTablename>'
    );
    
    INSERT INTO clickhouse_output
    SELECT
      id,
      name,
      age,
      rate
    FROM clickhouse_source;
  • Exemplo 2: Gravar em uma tabela distribuída.

    Considere uma tabela distribuída chamada distributed_table_test criada a partir de três tabelas locais chamadas local_table_test, localizadas nos nós 192.XX.XX.1, 192.XX.XX.2 e 192.XX.XX.3.

    • Para que o Flink grave dados diretamente nas tabelas locais e particione os dados por uma chave, use a seguinte DDL:

      CREATE TEMPORARY TABLE clickhouse_source (
        id INT,
        name VARCHAR,
        age BIGINT,
        rate FLOAT
      ) WITH (
        'connector' = 'datagen',
        'rows-per-second' = '50'
      );
      
      CREATE TEMPORARY TABLE clickhouse_output (
        id INT,
        name VARCHAR,
        age BIGINT,
        rate FLOAT
      ) WITH (
        'connector' = 'clickhouse',
        'url' = 'jdbc:clickhouse://192.XX.XX.1:3002,192.XX.XX.2:3002,192.XX.XX.3:3002/default',
        'userName' = '<yourUsername>',
        'password' = '<yourPassword>',
        'tableName' = 'local_table_test',
        'shardWrite' = 'true',
        'writeMode' = 'partition',
        'shardingKey' = 'name'
      );
      
      INSERT INTO clickhouse_output
      SELECT
        id,
        name,
        age,
        rate
      FROM clickhouse_source;
    • Se preferir que o Flink descubra automaticamente os nós da tabela local em vez de especificá-los manualmente na URL, utilize a seguinte DDL:

      CREATE TEMPORARY TABLE clickhouse_source (
        id INT,
        name VARCHAR,
        age BIGINT,
        rate FLOAT
      ) WITH (
        'connector' = 'datagen',
        'rows-per-second' = '50'
      );
      
      CREATE TEMPORARY TABLE clickhouse_output (
        id INT,
        name VARCHAR,
        age BIGINT,
        rate FLOAT
      ) WITH (
        'connector' = 'clickhouse',
        'url' = 'jdbc:clickhouse://192.XX.XX.1:3002/default', -- The JDBC URL of a node in the cluster.
        'userName' = '<yourUsername>',
        'password' = '<yourPassword>',
        'tableName' = 'distributed_table_test', -- The name of the distributed table.
        'shardWrite' = 'true',
        'inferLocalTable' = 'true', -- Set inferLocalTable to true.
        'writeMode' = 'partition',
        'shardingKey' = 'name'
      );
      
      INSERT INTO clickhouse_output
      SELECT
        id,
        name,
        age,
        rate
      FROM clickhouse_source;

Perguntas frequentes