Si necesitas ayuda, abre una incidencia en el repositorio o plantea tu pregunta en el Slack público de ClickHouse.
Licencia
Requisitos del entorno
Matriz de compatibilidad de versiones
Funcionalidades principales
- Incluye semántica exactly-once lista para usar. Se basa en una nueva funcionalidad del núcleo de ClickHouse llamada KeeperMap (utilizada como almacén de estado por el conector) y permite una arquitectura minimalista.
- Soporte para almacenes de estado de terceros: actualmente usa In-memory de forma predeterminada, pero puede utilizar KeeperMap (Redis se añadirá próximamente).
- Integración principal: desarrollada, mantenida y respaldada por ClickHouse.
- Probada continuamente con ClickHouse Cloud.
- Inserciones de datos con esquema declarado y sin esquema.
- Soporte para todos los tipos de datos de ClickHouse.
Instrucciones de instalación
Recopila los datos de conexión
Los detalles de su servicio de ClickHouse Cloud están disponibles en la consola de ClickHouse Cloud.
Seleccione un servicio y haga clic en Connect:
Elija HTTPS. Los detalles de conexión se muestran en un comando
curl de ejemplo.
Si usa ClickHouse autogestionado, los detalles de conexión los establece su administrador de ClickHouse.
Instrucciones generales de instalación
- Descargue un archivo ZIP que contenga el archivo JAR del conector desde la página de Releases del repositorio de ClickHouse Kafka Connect Sink.
- Extraiga el contenido del archivo ZIP y cópielo en la ubicación deseada.
- Agregue la ruta al directorio del plugin en la configuración plugin.path de su archivo de propiedades de Connect para que Confluent Platform pueda encontrar el plugin.
- Proporcione en la configuración un nombre de topic, el hostname de la instancia de ClickHouse y la contraseña.
- Reinicie Confluent Platform.
- Si usa Confluent Platform, inicie sesión en la UI de Confluent Control Center para verificar que ClickHouse Sink esté disponible en la lista de conectores disponibles.
Opciones de configuración
- detalles de conexión: hostname (obligatorio) y puerto (opcional)
- credenciales de usuario: contraseña (obligatoria) y nombre de usuario (opcional)
- clase del conector:
com.clickhouse.kafka.connect.ClickHouseSinkConnector(obligatoria) - topics o topics.regex: los topics de Kafka que se van a consultar; los nombres de los topics deben coincidir con los nombres de las tablas (obligatorio)
- convertidores de clave y valor: configúrelos en función del tipo de datos de su topic. Es obligatorio si aún no están definidos en la configuración del worker.
Tablas de destino
Preprocesamiento
Tipos de datos compatibles
-
(1) - JSON solo es compatible cuando la configuración de ClickHouse incluye
input_format_binary_read_json_as_string=1. Esto solo funciona con la familia de formatos RowBinary, y la configuración afecta a todas las columnas de la solicitud de insert, por lo que todas deben ser de tipo cadena. En este caso, el conector convertirá STRUCT en una cadena JSON. -
(2) - Cuando struct tiene unions como
oneof, el convertidor debe configurarse para NO añadir prefijos/sufijos a los nombres de los campos. Existe la configuracióngenerate.index.for.unions=falseparaProtobufConverter.
Recetas de configuración
Configuración básica
localhost:8443 con SSL habilitado; los datos están en JSON sin esquema.
La configuración del conector anterior requiere que habilites la sobrescritura de la configuración del cliente en la configuración de tu worker mediante
connector.client.config.override.policy=All. Consulta la documentación de Kafka Connect para obtener más información.Configuración básica con varios topics
Configuración básica con DLQ
Uso con distintos formatos de datos
Compatibilidad con esquemas Avro
Asignación de tipos de Avro
io.confluent.connect.avro.AvroConverter, la implementación oficial de serialización/deserialización de Avro en Kafka Connect. Consulta la documentación de Kafka Connect para obtener información avanzada sobre la lógica de conversión.
✅: Compatible
❌: No compatible
️⚠️: Parcialmente compatible
Consulta Tipos de datos compatibles para ver la asignación entre los tipos de Kafka Connect y los tipos de ClickHouse.
Esquemas de Avro no compatibles
- tipo lógico
decimalde longitud fija
- uniones nullable
- uniones de registro
Compatibilidad con esquemas de Protobuf
Asignación de tipos de Protobuf
io.confluent.connect.protobuf.ProtobufConverter, la implementación oficial de serialización/deserialización de Protobuf en Kafka Connect. Consulta la documentación de Kafka Connect para obtener información avanzada sobre la lógica de conversión.
✅: Compatible
❌: No compatible
️⚠️: Compatibilidad parcial
Consulta Tipos de datos compatibles para ver la asignación entre los tipos de Kafka Connect y los de ClickHouse.
Nota sobre la conversión de campos oneof en columnas de ClickHouse
oneof) al tipo Variant de ClickHouse. En su lugar, enumera los campos oneof como campos anulables individuales en el esquema de tu tabla de ClickHouse.
Por ejemplo:
Esquemas Protobuf no compatibles
- uniones de varios mensajes (antes de la versión 26.1 de CH)
allow_experimental_nullable_tuple_type=1 (consulta esta página de la documentación).
Compatibilidad con esquemas JSON
Compatibilidad con el convertidor String
Almacenamiento en búfer interno
poll() y los vuelque a ClickHouse en lotes más grandes. Esto puede mejorar el throughput en cargas de trabajo donde cada sondeo produce muchos lotes pequeños por partición.
Comportamiento clave:
bufferCountcontrola cuántos registros se almacenan en el búfer antes de volcarse.bufferFlushTimeestablece un tiempo máximo de espera (en milisegundos) antes de volcar los registros almacenados en el búfer.bufferFlushTimesolo surte efecto cuandobufferCount > 0.bufferCount=0ybufferFlushTime=0mantienen el almacenamiento en búfer deshabilitado (comportamiento predeterminado).- El almacenamiento en búfer no es compatible cuando
exactlyOnce=true.
exactlyOnce=false en la configuración de tu conector, o deshabilita el almacenamiento en búfer con bufferCount=0.
Ejemplo:
Registro
Monitoreo
Métricas específicas de ClickHouse
Métricas de productor/consumidor de Kafka
records-sent-total: Número total de registros enviados al topicbytes-sent-total: Número total de bytes enviados al topicrecord-send-rate: Tasa promedio de registros enviados por segundobyte-rate: Promedio de bytes enviados por segundocompression-rate: Ratio de compresión alcanzado
records-sent-total: Número total de registros enviados a la particiónbytes-sent-total: Número total de bytes enviados a la particiónrecords-lag: Lag actual de la particiónrecords-lead: Lead actual de la particiónreplica-fetch-lag: Información de lag de las réplicas
connection-creation-total: Número total de conexiones creadas con el nodo de Kafkaconnection-close-total: Número total de conexiones cerradasrequest-total: Número total de solicitudes enviadas al nodoresponse-total: Número total de respuestas recibidas del nodorequest-rate: Tasa promedio de solicitudes por segundoresponse-rate: Tasa promedio de respuestas por segundo
- Throughput: Hacer seguimiento de las tasas de ingestión de datos
- Lag: Identificar cuellos de botella y retrasos de procesamiento
- Compresión: Medir la eficiencia de la compresión de datos
- Estado de las conexiones: Supervisar la conectividad y la estabilidad de la red
Métricas de Kafka Connect Framework
task-count: Número total de tareas en el conectorrunning-task-count: Número de tareas que se están ejecutando actualmentepaused-task-count: Número de tareas actualmente en pausafailed-task-count: Número de tareas que han falladodestroyed-task-count: Número de tareas destruidasunassigned-task-count: Número de tareas sin asignar
running, paused, failed, destroyed, unassigned
Métricas de errores:
deadletterqueue-produce-failures: Número de escrituras fallidas en la DLQdeadletterqueue-produce-requests: Número total de intentos de escritura en la DLQlast-error-timestamp: Marca temporal del último errorrecords-skip-total: Número total de registros omitidos debido a erroresrecords-retry-total: Número total de registros reintentadoserrors-total: Número total de errores detectados
offset-commit-failures: Número de commits de offset fallidosoffset-commit-avg-time-ms: Tiempo promedio de los commits de offsetoffset-commit-max-time-ms: Tiempo máximo de los commits de offsetput-batch-avg-time-ms: Tiempo promedio para procesar un loteput-batch-max-time-ms: Tiempo máximo para procesar un lotesource-record-poll-total: Total de registros recuperados
Prácticas recomendadas de monitorización
- Supervise el consumer lag: Controle
records-lagpor partición para identificar cuellos de botella en el procesamiento - Controle las tasas de error: Observe
errors-totalyrecords-skip-totalpara detectar problemas de calidad de los datos - Observe el estado de las tareas: Supervise las métricas de estado de las tareas para asegurarse de que se ejecutan correctamente
- Mida el rendimiento: Use
records-send-rateybyte-ratepara controlar el rendimiento de la ingestión - Supervise el estado de las conexiones: Compruebe las métricas de conexión a nivel de nodo para detectar problemas de red
- Controle la eficiencia de la compresión: Use
compression-ratepara optimizar la transferencia de datos
Limitaciones
- No se admiten eliminaciones.
- El tamaño del lote se hereda de las propiedades del consumer de Kafka.
- Al usar KeeperMap para exactly-once, si se cambia o se restablece el offset, debe eliminar el contenido de KeeperMap para ese topic concreto. (Consulte la guía de solución de problemas a continuación para obtener más detalles)
Ajuste del rendimiento y optimización del caudal
¿Cuándo se necesita el ajuste del rendimiento?
- Cargas de trabajo de alto rendimiento: al procesar millones de eventos por segundo desde topics de Kafka
- Consumer lag: cuando el conector no puede seguir el ritmo de producción de datos, lo que provoca un retraso cada vez mayor
- Restricciones de recursos: cuando necesitas optimizar el uso de CPU, memoria o red
- Múltiples topics: al consumir simultáneamente de varios topics de gran volumen
- Mensajes pequeños: al trabajar con muchos mensajes pequeños que se beneficiarían del procesamiento por lotes en el servidor
- Procesas volúmenes bajos o moderados (< 10,000 mensajes/segundo)
- El consumer lag es estable y aceptable para tu caso de uso
- La configuración predeterminada del conector ya cumple tus requisitos de rendimiento
- Tu clúster de ClickHouse puede manejar fácilmente la carga entrante
Comprender el flujo de datos
- Kafka Connect Framework recupera mensajes de los topics de Kafka en segundo plano
- El conector sondea mensajes del búfer interno del framework
- El conector agrupa los mensajes en lotes según el tamaño del sondeo
- ClickHouse recibe el insert por lotes mediante HTTP/S
- ClickHouse procesa el insert (de forma síncrona o asíncrona)
Ajuste del tamaño del lote en Kafka Connect
Configuración de fetch
fetch.min.bytes: Cantidad mínima de datos antes de que el framework entregue datos al conector (predeterminado: 1 byte)fetch.max.bytes: Cantidad máxima de datos que se puede recuperar en una sola solicitud (predeterminado: 52428800 / 50 MB)fetch.max.wait.ms: Tiempo máximo de espera antes de devolver datos si no se alcanzafetch.min.bytes(predeterminado: 500 ms)
En Confluent Cloud, para ajustar esta configuración es necesario abrir un caso de soporte a través de Confluent Cloud.
Configuración de sondeo
max.poll.records: Número máximo de registros devueltos en un único sondeo (predeterminado: 500)max.partition.fetch.bytes: Cantidad máxima de datos por partición (predeterminado: 1048576 / 1 MB)
En Confluent Cloud, para ajustar estos parámetros, es necesario abrir un caso de soporte a través de Confluent Cloud.
Configuración recomendada para un alto rendimiento
Las propiedades anteriores requieren que habilites las anulaciones del cliente en la configuración de tu worker mediante
connector.client.config.override.policy=All. Consulta la documentación de Kafka Connect para obtener más información.- Lotes más grandes = Mejor rendimiento de ingestión en ClickHouse, menos partes y menor sobrecarga
- Lotes más grandes = Mayor uso de memoria y posible aumento de la latencia de extremo a extremo
- Lotes demasiado grandes = Riesgo de timeouts, errores OutOfMemory o de exceder
max.poll.interval.ms
Inserciones asíncronas
Cuándo usar inserción asíncrona
- Muchos lotes pequeños: Tu conector envía lotes pequeños con frecuencia (< 1000 filas por lote)
- Alta concurrencia: Varias tareas del conector escriben en la misma tabla
- Implementación distribuida: Ejecutas muchas instancias del conector en distintos hosts
- Sobrecarga por creación de partes: Estás experimentando errores de “too many parts”
- Carga de trabajo mixta: Combinas la ingestión en tiempo real con cargas de trabajo de consultas
- Ya estás enviando lotes grandes (> 10,000 filas por lote) con una frecuencia controlada
- Necesitas visibilidad inmediata de los datos (las consultas deben ver los datos al instante)
- La semántica exactly-once con
wait_for_async_insert=0entra en conflicto con tus requisitos - Tu caso de uso puede beneficiarse más de mejoras en la agrupación en lotes del lado del cliente
Cómo funcionan las inserciones asíncronas
- Recibe la consulta de inserción del conector
- Escribe los datos en un búfer en memoria (en lugar de escribirlos inmediatamente en disco)
- Devuelve una respuesta de éxito al conector (si
wait_for_async_insert=0) - Vuelca el búfer a disco cuando se cumple una de estas condiciones:
- El búfer alcanza
async_insert_max_data_size(valor predeterminado: 100 MB) - Han transcurrido
async_insert_busy_timeout_msmilisegundos desde la primera inserción (valor predeterminado: 1000 ms) - Se alcanza el número máximo de consultas acumuladas (
async_insert_max_query_number, valor predeterminado: 100)
- El búfer alcanza
Habilitar inserción asíncrona
async insert al parámetro de configuración clickhouseSettings:
async_insert=1: Habilita las inserciones asíncronaswait_for_async_insert=1(recomendado): El conector espera a que los datos se escriban en el almacenamiento de ClickHouse antes de confirmar la recepción. Proporciona garantías de entrega.wait_for_async_insert=0: El conector confirma la recepción inmediatamente después de almacenarlos en búfer. Ofrece mejor rendimiento, pero los datos pueden perderse si el servidor falla antes de que se escriban en disco.
Ajuste del comportamiento de async insert
async insert:
async_insert_max_data_size(valor predeterminado: 104857600 / 100 MB): Tamaño máximo del búfer antes del vaciadoasync_insert_busy_timeout_ms(valor predeterminado: 1000): Tiempo máximo (ms) antes del vaciadoasync_insert_stale_timeout_ms(valor predeterminado: 0): Tiempo (ms) desde la última inserción antes del vaciadoasync_insert_max_query_number(valor predeterminado: 100): Número máximo de consultas antes del vaciado
- Beneficios: Menos partes, mejor rendimiento de las fusiones, menor sobrecarga de CPU, mayor rendimiento con alta concurrencia
- Consideraciones: Los datos no se pueden consultar de inmediato, latencia de extremo a extremo ligeramente mayor
- Riesgos: Pérdida de datos si el servidor falla y
wait_for_async_insert=0, posible presión de memoria con búferes grandes
Inserciones asíncronas con semántica de exactamente una vez
exactlyOnce=true con inserciones asíncronas:
wait_for_async_insert=1 con exactly-once para garantizar que las confirmaciones de offsets solo se produzcan después de que los datos se hayan guardado de forma persistente.
Para obtener más información sobre la inserción asíncrona, consulta la documentación de inserción asíncrona de ClickHouse.
Paralelismo del conector
Tareas por conector
- El número máximo de tareas efectivas = número de particiones del topic
- Cada tarea mantiene su propia conexión a ClickHouse
- Más tareas = más sobrecarga y posible contención de recursos
tasks.max igual al número de particiones del topic y luego ajústalo en función de la CPU y las métricas de caudal.
Ignorar las particiones al agrupar en lotes
exactlyOnce=false. Esta configuración puede mejorar el caudal al crear lotes más grandes, pero se pierden las garantías de orden dentro de cada partición.
Varios topics de alto caudal
topic2TableMap para asignar topics a tablas y se está produciendo un cuello de botella en la inserción que provoca consumer lag, considere crear un conector por topic en su lugar.
La razón principal es que, actualmente, los lotes se insertan en cada tabla en serie.
Recomendación: Para varios topics de alto volumen, implemente una instancia de conector por topic para maximizar el caudal de inserción en paralelo.
Consideraciones sobre el motor de tabla de ClickHouse
MergeTree: La mejor opción para la mayoría de los casos de uso; equilibra el rendimiento de las consultas y las insercionesReplicatedMergeTree: Requerido para alta disponibilidad; añade sobrecarga de replicación*MergeTreecon unORDER BYadecuado: Optimiza según tus patrones de consulta
Pool de conexiones y tiempos de espera
socket_timeout(predeterminado: 30000 ms): Tiempo máximo para las operaciones de lecturaconnection_timeout(predeterminado: 10000 ms): Tiempo máximo para establecer la conexión
Supervisión y solución de problemas de rendimiento
- Consumer lag: Use las herramientas de supervisión de Kafka para seguir el retraso por partición
- Métricas del conector: Supervise
receivedRecords,recordProcessingTime,taskProcessingTimemediante JMX (consulte Monitoring) - Métricas de ClickHouse:
system.asynchronous_inserts: Supervise el uso del búfer de inserción asíncronasystem.parts: Supervise el número de partes para detectar problemas de fusionessystem.merges: Supervise las fusiones activassystem.events: SigaInsertedRows,InsertedBytes,FailedInsertQuery
Resumen de prácticas recomendadas
- Empieza con los valores predeterminados y luego mide y ajusta en función del rendimiento real
- Prefiere lotes más grandes: apunta a 10,000-100,000 filas por inserción cuando sea posible
- Usa inserción asíncrona cuando envíes muchos lotes pequeños o haya alta concurrencia
- Usa siempre
wait_for_async_insert=1con semántica exactly-once - Escala horizontalmente: aumenta
tasks.maxhasta el número de particiones - Un conector por topic de alto volumen para obtener el máximo caudal
- Supervisa continuamente: controla el consumer lag, la cantidad de partes y la actividad de fusión
- Haz pruebas exhaustivas: prueba siempre los cambios de configuración con una carga realista antes de la implementación en producción
Ejemplo: Configuración de alto caudal
La configuración del conector anterior requiere que habilites las anulaciones de cliente en la configuración de tu worker mediante
connector.client.config.override.policy=All. Consulta la documentación de Kafka Connect para obtener más información.- Procesa hasta 10.000 registros por ciclo de sondeo
- Agrupa en lotes entre particiones para inserciones más grandes
- Usa inserción asíncrona con un búfer de 16 MB
- Ejecuta 8 tareas en paralelo (haz que coincida con tu número de particiones)
- Está optimizada para el caudal por encima del orden estricto
Solución de problemas
”Incongruencia de estado para el topic [someTopic] partición [0]”
Este ajuste puede afectar las garantías exactly-once.
”¿Qué errores reintentará el conector?”
ClickHouseException- Esta es una excepción genérica que puede lanzar ClickHouse. Suele producirse cuando el server está sobrecargado, y los siguientes códigos de error se consideran especialmente transitorios:- 3 - UNEXPECTED_END_OF_FILE
- 107 - FILE_DOESNT_EXIST
- 159 - TIMEOUT_EXCEEDED
- 164 - READONLY
- 202 - TOO_MANY_SIMULTANEOUS_QUERIES
- 203 - NO_FREE_CONNECTION
- 209 - SOCKET_TIMEOUT
- 210 - NETWORK_ERROR
- 241 - MEMORY_LIMIT_EXCEEDED
- 242 - TABLE_IS_READ_ONLY
- 252 - TOO_MANY_PARTS
- 285 - TOO_FEW_LIVE_REPLICAS
- 319 - UNKNOWN_STATUS_OF_INSERT
- 425 - SYSTEM_ERROR
- 999 - KEEPER_EXCEPTION
SocketTimeoutException- Se lanza cuando el socket agota el tiempo de espera.UnknownHostException- Se lanza cuando no se puede resolver el host.IOException- Se lanza cuando hay un problema con la red.
”Todos mis datos están en blanco o en cero”
_ como delimitador). Los campos de la tabla seguirán entonces el formato “field1_field2_field3” (es decir, “before_id”, “after_id”, etc.).
”Quiero usar las claves de Kafka en ClickHouse”
value, pero puedes usar la transformación KeyToValue para mover la clave al campo value (con el nuevo nombre de campo _key):