Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Conector Flink

Suportado pelo ClickHouse

Este é o Conector Sink oficial do Apache Flink, com suporte da ClickHouse. Ele foi desenvolvido com o AsyncSinkBase do Flink e o Java client oficial do ClickHouse.

O conector oferece suporte à API do DataStream do Apache Flink. O suporte à Table API está planejado para um lançamento futuro.

Requisitos

  • Java 11+ (para o Flink 1.17+) ou 17+ (para o Flink 2.0+)
  • Apache Flink 1.17+

O conector foi dividido em dois artefatos para oferecer suporte ao Flink 1.17+ e ao Flink 2.0+. Escolha o artefato correspondente à versão do Flink que você deseja usar:

Flink Version Artifact ClickHouse Java Client Version Required Java
latest flink-connector-clickhouse-2.0.0 0.9.5 Java 17+
2.0.1 flink-connector-clickhouse-2.0.0 0.9.5 Java 17+
2.0.0 flink-connector-clickhouse-2.0.0 0.9.5 Java 17+
1.20.2 flink-connector-clickhouse-1.17 0.9.5 Java 11+
1.19.3 flink-connector-clickhouse-1.17 0.9.5 Java 11+
1.18.1 flink-connector-clickhouse-1.17 0.9.5 Java 11+
1.17.2 flink-connector-clickhouse-1.17 0.9.5 Java 11+

Instalação e configuração

Importar como dependência

<dependency>
    <groupId>com.clickhouse.flink</groupId>
    <artifactId>flink-connector-clickhouse-2.0.0</artifactId>
    <version>{{ stable_version }}</version>
    <classifier>all</classifier>
</dependency>
<dependency>
    <groupId>com.clickhouse.flink</groupId>
    <artifactId>flink-connector-clickhouse-1.17</artifactId>
    <version>{{ stable_version }}</version>
    <classifier>all</classifier>
</dependency>

Baixe o binário

O padrão de nomenclatura do arquivo JAR binário é:

flink-connector-clickhouse-${flink_version}-${stable_version}-all.jar

onde:

Você pode encontrar todos os arquivos JAR disponíveis já lançados no Maven Central Repository.

Como usar a API do DataStream

Trecho

Digamos que você queira inserir dados CSV brutos no ClickHouse:

public static void main(String[] args) {
    // Configure o ClickHouseClient
    ClickHouseClientConfig clientConfig = new ClickHouseClientConfig(url, username, password, database, tableName);

    // Crie um ElementConverter
    ElementConverter<String, ClickHousePayload> convertorString = new ClickHouseConvertor<>(String.class);

    // Crie o sink e defina o formato usando `setClickHouseFormat`
    ClickHouseAsyncSink<String> csvSink = new ClickHouseAsyncSink<>(
            convertorString,
            MAX_BATCH_SIZE,
            MAX_IN_FLIGHT_REQUESTS,
            MAX_BUFFERED_REQUESTS,
            MAX_BATCH_SIZE_IN_BYTES,
            MAX_TIME_IN_BUFFER_MS,
            MAX_RECORD_SIZE_IN_BYTES,
            clientConfig
    );

    csvSink.setClickHouseFormat(ClickHouseFormat.CSV);

    // Por fim, conecte seu DataStream ao sink.
    final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    Path csvFilePath = new Path(fileFullName);
    FileSource<String> csvSource = FileSource
            .forRecordStreamFormat(new TextLineInputFormat(), csvFilePath)
            .build();

    env.fromSource(
            csvSource,
            WatermarkStrategy.noWatermarks(),
            "GzipCsvSource"
    ).sinkTo(csvSink);
}

Mais exemplos e trechos de código podem ser encontrados em nossos testes:

Exemplo de início rápido

Criamos um exemplo baseado em Maven para facilitar os primeiros passos com o ClickHouse Sink:

Para instruções mais detalhadas, consulte o Guia de exemplos

Opções de conexão com a API DataStream

Opções do cliente ClickHouse

Parâmetros Descrição Valor padrão Obrigatório
url URL completa do ClickHouse N/A Sim
username Nome de usuário do banco de dados ClickHouse N/A Sim
password Senha do banco de dados ClickHouse N/A Sim
database Nome do banco de dados ClickHouse N/A Sim
table Nome da tabela ClickHouse N/A Sim
options Mapa de opções de configuração do Java client Mapa vazio Não
serverSettings Mapa de configurações de sessão do servidor ClickHouse Mapa vazio Não
enableJsonSupportAsString Configuração do servidor ClickHouse para esperar uma String formatada em JSON para o tipo de dado JSON true Não

options e serverSettings devem ser passados ao cliente como Map<String, String>. Um mapa vazio em qualquer um deles usará os padrões do cliente ou do servidor, respectivamente.

Por exemplo:

Map<String, String> javaClientOptions = Map.of(
    ClientConfigProperties.CA_CERTIFICATE.getKey(), "<my_CA_cert>",
    ClientConfigProperties.SSL_CERTIFICATE.getKey(), "<my_SSL_cert>",
    ClientConfigProperties.CLIENT_NETWORK_BUFFER_SIZE.getKey(), "30000",
    ClientConfigProperties.HTTP_MAX_OPEN_CONNECTIONS.getKey(), "5"
);

Map<String, String> serverSettings = Map.of(
    "insert_deduplicate", "1"
);

ClickHouseClientConfig clientConfig = new ClickHouseClientConfig(
    url,
    username,
    password,
    database,
    tableName,
    javaClientOptions,
    serverSettings,
    false // enableJsonSupportAsString
);

Opções do sink

As opções a seguir vêm diretamente do AsyncSinkBase do Flink:

Parâmetros Descrição Valor padrão Obrigatório
maxBatchSize Número máximo de registros inseridos em um único lote N/A Sim
maxInFlightRequests Número máximo de solicitações em andamento permitido antes de o sink aplicar backpressure N/A Sim
maxBufferedRequests Número máximo de registros que podem ficar em buffer no sink antes de o backpressure ser aplicado N/A Sim
maxBatchSizeInBytes Tamanho máximo (em bytes) que um lote pode atingir. Todos os lotes enviados serão menores ou iguais a esse tamanho N/A Sim
maxTimeInBufferMS Tempo máximo que um registro pode permanecer no sink antes de ser gravado N/A Sim
maxRecordSizeInBytes Tamanho máximo de registro que o sink aceitará; registros maiores que isso serão rejeitados automaticamente N/A Sim

Tipos de dados compatíveis

A tabela abaixo traz uma referência rápida para a conversão de tipos de dados ao inserir dados do Flink no ClickHouse.

Tipo Java Tipo do ClickHouse Suportado Método de serialização
byte/Byte Int8 DataWriter.writeInt8
short/Short Int16 DataWriter.writeInt16
int/Integer Int32 DataWriter.writeInt32
long/Long Int64 DataWriter.writeInt64
BigInteger Int128 DataWriter.writeInt128
BigInteger Int256 DataWriter.writeInt256
short/Short UInt8 DataWriter.writeUInt8
int/Integer UInt8 DataWriter.writeUInt8
int/Integer UInt16 DataWriter.writeUInt16
long/Long UInt32 DataWriter.writeUInt32
long/Long UInt64 DataWriter.writeUInt64
BigInteger UInt64 DataWriter.writeUInt64
BigInteger UInt128 DataWriter.writeUInt128
BigInteger UInt256 DataWriter.writeUInt256
BigDecimal Decimal DataWriter.writeDecimal
BigDecimal Decimal32 DataWriter.writeDecimal
BigDecimal Decimal64 DataWriter.writeDecimal
BigDecimal Decimal128 DataWriter.writeDecimal
BigDecimal Decimal256 DataWriter.writeDecimal
float/Float Float DataWriter.writeFloat32
double/Double Double DataWriter.writeFloat64
boolean/Boolean Boolean DataWriter.writeBoolean
String String DataWriter.writeString
String FixedString DataWriter.writeFixedString
LocalDate Date DataWriter.writeDate
LocalDate Date32 DataWriter.writeDate32
LocalDateTime DateTime DataWriter.writeDateTime
ZonedDateTime DateTime DataWriter.writeDateTime
LocalDateTime DateTime64 DataWriter.writeDateTime64
ZonedDateTime DateTime64 DataWriter.writeDateTime64
int/Integer Time N/A
long/Long Time64 N/A
byte/Byte Enum8 DataWriter.writeInt8
int/Integer Enum16 DataWriter.writeInt16
java.util.UUID UUID DataWriter.writeIntUUID
String JSON DataWriter.writeJSON
Array<Type> Array<Type> DataWriter.writeArray
Map<K,V> Map<K,V> DataWriter.writeMap
Tuple<Type,..> Tuple<T1,T2,..> DataWriter.writeTuple
Object Variant N/A

Observações:

  • Um ZoneId deve ser fornecido ao realizar operações com data.
  • Precisão e escala devem ser fornecidas ao realizar operações decimais.
  • Para que o ClickHouse consiga interpretar uma String Java como JSON, é necessário habilitar enableJsonSupportAsString em ClickHouseClientConfig.
  • O conector requer um ElementConvertor para mapear elementos no DataStream de entrada para payloads do ClickHouse. Para isso, o conector fornece ClickHouseConvertor e POJOConvertor, que podem ser usados para implementar esse mapeamento com os métodos de serialização de DataWriter acima.

Formatos de entrada suportados

Você pode encontrar a lista de formatos de entrada disponíveis do ClickHouse nesta página da documentação e em ClickHouseFormat.java.

Para especificar o formato que o conector deve usar para serializar seu DataStream como payloads para o ClickHouse, use a função setClickHouseFormat. Por exemplo:

ClickHouseAsyncSink<String> csvSink = new ClickHouseAsyncSink<>(
        convertorString,
        MAX_BATCH_SIZE,
        MAX_IN_FLIGHT_REQUESTS,
        MAX_BUFFERED_REQUESTS,
        MAX_BATCH_SIZE_IN_BYTES,
        MAX_TIME_IN_BUFFER_MS,
        MAX_RECORD_SIZE_IN_BYTES,
        clientConfig
);
csvSink.setClickHouseFormat(ClickHouseFormat.CSV);

Métricas

O conector expõe as seguintes métricas adicionais, além das métricas já existentes do Flink:

Métrica Descrição Tipo Status
numBytesSend Número total de bytes enviados ao ClickHouse no payload da requisição. Observação: esta métrica mede o tamanho dos dados serializados enviados pela rede e pode diferir de written_bytes do ClickHouse em system.query_log, que reflete os bytes efetivamente gravados no armazenamento após o processamento Contador
numRecordSend Número total de registros enviados ao ClickHouse Contador
numRequestSubmitted Número total de requisições enviadas (número real de flushes executados) Contador
numOfDroppedBatches Número total de lotes descartados devido a falhas não recuperáveis Contador
numOfDroppedRecords Número total de registros descartados devido a falhas não recuperáveis Contador
totalBatchRetries Número total de novas tentativas de lote devido a falhas recuperáveis Contador
writeLatencyHistogram Histograma da distribuição da latência de gravações bem-sucedidas (ms) Histograma
writeFailureLatencyHistogram Histograma da distribuição da latência de gravações com falha (ms) Histograma
triggeredByMaxBatchSizeCounter Número total de flushes acionados ao atingir maxBatchSize Contador
triggeredByMaxBatchSizeInBytesCounter Número total de flushes acionados ao atingir maxBatchSizeInBytes Contador
triggeredByMaxTimeInBufferMSCounter Número total de flushes acionados ao atingir maxTimeInBufferMS Contador
actualRecordsPerBatch Histograma da distribuição do tamanho real do lote Histograma
actualBytesPerBatch Histograma da distribuição real de bytes por lote Histograma

Limitações

  • No momento, o sink oferece uma garantia de entrega at-least-once. O suporte à semântica exactly-once está sendo acompanhado aqui.
  • O sink ainda não oferece suporte a uma fila de dead-letter (DLQ) para armazenar temporariamente registros que não podem ser processados. Enquanto isso, o conector tentará reinserir os registros com falha e os descartará em caso de insucesso. Esse recurso está sendo acompanhado aqui.
  • O sink ainda não oferece suporte à criação por meio da Table API do Flink ou do Flink SQL. Esse recurso está sendo acompanhado aqui.

Compatibilidade de versões do ClickHouse e segurança

  • O conector é testado diariamente, por meio de um workflow de CI, com uma variedade de versões recentes do ClickHouse, incluindo latest e head. As versões testadas são atualizadas periodicamente à medida que novos lançamentos do ClickHouse entram em atividade. Veja aqui as versões com as quais o conector é testado diariamente.
  • Consulte a política de segurança do ClickHouse para ver vulnerabilidades de segurança conhecidas e como relatar uma vulnerabilidade.
  • Recomendamos atualizar o conector continuamente para não perder correções de segurança e outras melhorias.
  • Se você tiver algum problema com a migração, crie uma issue no GitHub e responderemos!
  • Para obter o melhor desempenho, garanta que o tipo de elemento do seu DataStream não seja um tipo genérico — veja aqui a distinção de tipos do Flink. Elementos não genéricos evitam a sobrecarga de serialização do Kryo e melhoram a vazão para o ClickHouse.
  • Recomendamos definir maxBatchSize para pelo menos 1000 e, idealmente, entre 10.000 e 100.000. Veja este guia sobre inserções em massa para mais informações.
  • Para fazer desduplicação no estilo OLTP ou upsert no ClickHouse, consulte esta página da documentação. Observação: isso não deve ser confundido com a desduplicação em lote que ocorre em novas tentativas.

Solução de problemas

CANNOT_READ_ALL_DATA

O erro a seguir pode ocorrer:

com.clickhouse.client.api.ServerException: Code: 33. DB::Exception: Cannot read all data. Bytes read: 9205. Bytes expected: 1100022.: (at row 9) : While executing BinaryRowInputFormat. (CANNOT_READ_ALL_DATA)

Causa: Na maioria dos casos, o erro CANNOT_READ_ALL_DATA significa que o schema da sua table do ClickHouse divergiu do schema do registro no Flink. Isso pode acontecer quando um deles é alterado de uma forma incompatível com versões anteriores.

Solução: Atualize o schema da sua table do ClickHouse ou o tipo de dado de entrada do conector (ou ambos) para que sejam compatíveis. Se necessário, consulte o mapeamento de tipos para ver como mapear tipos Java para tipos do ClickHouse. Observação: se ainda houver registros em trânsito, você precisará redefinir o state do Flink ao reiniciar o conector.

Baixa vazão

Você pode notar que a vazão do conector não escala com o paralelismo do job (número de tasks do Flink) ao gravar no ClickHouse.

Causa: o processo de merge de parts em segundo plano do ClickHouse pode estar reduzindo a velocidade das inserções. Isso pode acontecer quando o tamanho de lote configurado é muito pequeno, o conector está fazendo flush com muita frequência, ou por uma combinação dos dois fatores.

Solução: monitore as métricas numRequestSubmitted e actualRecordsPerBatch para ajudar a determinar como ajustar o tamanho do lote (maxBatchSize) e a frequência de flush. Além disso, consulte Uso avançado e recomendado para recomendações de dimensionamento de lote.

Faltam linhas na minha tabela do ClickHouse

Causa: O(s) lote(s) foi(foram) descartado(s) devido a uma falha não recuperável ou porque não pôde(ram) ser inserido(s) dentro do número configurado de tentativas (configurável via ClickHouseClientConfig.setNumberOfRetries()). Observação: por padrão, o conector tentará reinserir um lote em até 3 tentativas antes de descartá-lo.

Solução: Inspecione os logs do TaskManager e/ou os stack traces para identificar a causa raiz.

Contribuição e suporte

Se você quiser contribuir com o projeto ou relatar algum problema, sua colaboração será muito bem-vinda! Visite nosso repositório no GitHub para abrir uma issue, sugerir melhorias ou enviar um pull request.

Contribuições são bem-vindas! Consulte o guia de contribuição no repositório antes de começar. Obrigado por ajudar a melhorar o conector do ClickHouse para Flink!

Navigation