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+
Matriz de compatibilidade das versões do Flink
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
Para o Flink 2.0+
<dependency>
<groupId>com.clickhouse.flink</groupId>
<artifactId>flink-connector-clickhouse-2.0.0</artifactId>
<version>{{ stable_version }}</version>
<classifier>all</classifier>
</dependency>dependencies {
implementation("com.clickhouse.flink:flink-connector-clickhouse-2.0.0:{{ stable_version }}")
}libraryDependencies += "com.clickhouse.flink" % "flink-connector-clickhouse-2.0.0" % {{ stable_version }} classifier "all"Para o Flink 1.17+
<dependency>
<groupId>com.clickhouse.flink</groupId>
<artifactId>flink-connector-clickhouse-1.17</artifactId>
<version>{{ stable_version }}</version>
<classifier>all</classifier>
</dependency>dependencies {
implementation("com.clickhouse.flink:flink-connector-clickhouse-1.17:{{ stable_version }}")
}libraryDependencies += "com.clickhouse.flink" % "flink-connector-clickhouse-1.17" % {{ stable_version }} classifier "all"Baixe o binário
O padrão de nomenclatura do arquivo JAR binário é:
flink-connector-clickhouse-${flink_version}-${stable_version}-all.jaronde:
flink_versioné2.0.0ou1.17stable_versioné uma versão estável do artefato
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.
Inserção de 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
ZoneIddeve 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
enableJsonSupportAsStringemClickHouseClientConfig. - O conector requer um
ElementConvertorpara mapear elementos noDataStreamde entrada para payloads do ClickHouse. Para isso, o conector forneceClickHouseConvertorePOJOConvertor, que podem ser usados para implementar esse mapeamento com os métodos de serialização deDataWriteracima.
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!
Uso avançado e recomendado
- 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
maxBatchSizepara 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!