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 |
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 |
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.
NotaPara 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.parallelismnã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
ignoreDeleteestá 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 é Nota
Para gravar dados em uma tabela distribuída do ClickHouse, a |
|
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:
Nota
Se você definir |
|
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:
|
|
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:
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:
Nota
Se você definir |
|
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:
Nota
|
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 |
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>
);
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;
-