Saltar al contenido principal
El motor de tabla Kafka se puede utilizar para leer datos de y escribir datos en Apache Kafka y otros brokers compatibles con la API de Kafka (p. ej., Redpanda, Amazon MSK).

De Kafka a ClickHouse

Si usas ClickHouse Cloud, te recomendamos usar ClickPipes en su lugar. ClickPipes admite de forma nativa conexiones de red privadas, el escalado independiente de los recursos de ingestión y del cluster, y una monitorización exhaustiva para transmitir datos de Kafka a ClickHouse.
Para usar el motor de tabla Kafka, deberías tener conocimientos generales sobre las vistas materializadas de ClickHouse.

Resumen

Inicialmente, nos centramos en el caso de uso más común: utilizar el motor de tabla de Kafka para insertar datos en ClickHouse desde Kafka. El motor de tabla de Kafka permite a ClickHouse leer directamente de un Kafka topic. Aunque es útil para ver mensajes en un topic, este motor, por diseño, solo permite una única lectura; es decir, cuando se ejecuta una consulta sobre la tabla, consume datos de la cola y avanza el offset del consumer antes de devolver los resultados al cliente. En la práctica, los datos no pueden volver a leerse sin restablecer esos offsets. Para conservar estos datos a partir de una lectura del motor de tabla, necesitamos un mecanismo para capturarlos e insertarlos en otra tabla. Las vistas materializadas basadas en triggers proporcionan esta funcionalidad de forma nativa. Una vista materializada inicia una lectura del motor de tabla y recibe lotes de documentos. La cláusula TO determina el destino de los datos, normalmente una tabla de la familia MergeTree. Este proceso se visualiza a continuación:

Pasos

1. Preparar
Si ya tiene datos en un topic de destino, puede adaptar lo siguiente para usarlo con su conjunto de datos. Como alternativa, se proporciona aquí un conjunto de datos de ejemplo de GitHub. Este conjunto de datos se utiliza en los ejemplos que aparecen a continuación e incluye un esquema reducido y un subconjunto de las filas (en concreto, nos limitamos a eventos de GitHub relacionados con el repositorio de ClickHouse), en comparación con el conjunto de datos completo disponible aquí, para simplificar. Aun así, sigue siendo suficiente para que funcionen la mayoría de las consultas publicadas con el conjunto de datos.
2. Configurar ClickHouse
Este paso es necesario si te conectas a un Kafka protegido. Estos ajustes no se pueden pasar mediante comandos SQL DDL y deben configurarse en el archivo config.xml de ClickHouse. Suponemos que te conectas a una instancia protegida con SASL. Este es el método más sencillo al interactuar con Confluent Cloud.
Coloque el fragmento anterior en un archivo nuevo dentro de su directorio conf.d/ o incorpórelo a los archivos de configuración existentes. Para ver los ajustes que se pueden configurar, consulte aquí. También vamos a crear una base de datos llamada KafkaEngine para utilizarla en este tutorial:
Una vez que hayas creado la base de datos, tendrás que cambiar a ella:
3. Cree la tabla de destino
Prepare la tabla de destino. En el ejemplo siguiente usamos el esquema reducido de GitHub por brevedad. Tenga en cuenta que, aunque usamos un motor de tabla MergeTree, este ejemplo puede adaptarse fácilmente a cualquier miembro de la familia MergeTree.
4. Crear y rellenar el topic
A continuación, vamos a crear un topic. Hay varias herramientas que podemos usar para hacerlo. Si estamos ejecutando Kafka localmente en nuestra máquina o dentro de un contenedor de Docker, RPK funciona bien. Podemos crear un topic llamado github con 5 particiones ejecutando el siguiente comando:
Si ejecutamos Kafka en Confluent Cloud, puede que prefiramos usar la CLI de Confluent:
Ahora necesitamos llenar este topic con algunos datos, lo que haremos con kcat. Podemos ejecutar un comando similar al siguiente si estamos ejecutando Kafka localmente con la autenticación deshabilitada:
O bien, lo siguiente si nuestro clúster de Kafka usa SASL para autenticarse:
El conjunto de datos contiene 200.000 filas, por lo que debería ingestarse en apenas unos segundos. Si quieres trabajar con un conjunto de datos más grande, consulta la sección sobre conjuntos de datos grandes del repositorio de GitHub de ClickHouse/kafka-samples.
5. Crear el motor de tabla de Kafka
El siguiente ejemplo crea un motor de tabla de Kafka con el mismo esquema que la tabla MergeTree. Esto no es estrictamente necesario, ya que puede tener columnas alias o efímeras en la tabla de destino. La configuración es importante; sin embargo, tenga en cuenta el uso de JSONEachRow como tipo de dato para consumir JSON desde un topic de Kafka. Los valores github y clickhouse representan el nombre del topic y los nombres de los grupos de consumidores, respectivamente. En realidad, los topics pueden ser una lista de valores.
A continuación, analizamos la configuración del motor y la optimización del rendimiento. En este punto, una consulta SELECT sencilla sobre la tabla github_queue debería leer algunas filas. Ten en cuenta que esto hará avanzar los offsets del consumidor, lo que impedirá volver a leer estas filas sin un reinicio. Fíjate en el límite y en el parámetro obligatorio stream_like_engine_allow_direct_select.
6. Crear la vista materializada
La vista materializada conectará las dos tablas creadas previamente, leerá datos desde el motor de tabla de Kafka y los insertará en la tabla MergeTree de destino. Podemos realizar varias transformaciones de datos. Haremos una lectura e inserción simples. El uso de * supone que los nombres de las columnas son idénticos (distinguen entre mayúsculas y minúsculas).
En el momento de su creación, la vista materializada se conecta al motor de Kafka y comienza a leer, insertando filas en la tabla de destino. Este proceso continuará indefinidamente, consumiendo las inserciones de mensajes posteriores en Kafka. Si lo deseas, puedes volver a ejecutar el script de inserción para añadir más mensajes a Kafka.
7. Confirme que se hayan insertado las filas
Confirme que haya datos en la tabla de destino:
Deberías ver 200.000 filas:

Operaciones comunes

Detener & reiniciar el consumo de mensajes
Para detener el consumo de mensajes, puede desvincular la tabla del motor Kafka:
Esto no afectará los offsets del grupo de consumidores. Para reanudar el consumo y continuar desde el offset anterior, vuelva a adjuntar la tabla.
Añadir metadatos de Kafka
Puede resultar útil hacer un seguimiento de los metadatos de los mensajes originales de Kafka una vez ingeridos en ClickHouse. Por ejemplo, puede que queramos saber cuánto hemos consumido de un topic o una partición concretos. Para ello, el motor de tabla de Kafka expone varias columnas virtuales. Estas pueden persistirse como columnas en nuestra tabla de destino modificando nuestro esquema y la instrucción SELECT de la vista materializada. Primero, realizamos la operación de detención descrita anteriormente antes de añadir columnas a nuestra tabla de destino.
A continuación, añadimos columnas informativas para identificar el topic de origen y la partición de la que procede la fila.
A continuación, debemos asegurarnos de que las columnas virtuales estén mapeadas según corresponda. Las columnas virtuales llevan el prefijo _. Puede consultar una lista completa de las columnas virtuales aquí. Para actualizar nuestra tabla con las columnas virtuales, tendremos que eliminar la vista materializada, volver a adjuntar la tabla con motor Kafka y volver a crear la vista materializada.
Las filas consumidas recientemente deberían incluir los metadatos.
El resultado es así:
Modificar la configuración del motor Kafka
Recomendamos eliminar la tabla con motor Kafka y volver a crearla con la nueva configuración. No es necesario modificar la vista materializada durante este proceso; el consumo de mensajes se reanudará una vez recreada la tabla con motor Kafka.
Depuración de incidencias
Errores como los problemas de autenticación no se notifican en las respuestas al DDL del motor Kafka. Para diagnosticar estos problemas, recomendamos usar el archivo de registro principal de ClickHouse, clickhouse-server.err.log. También se puede habilitar un registro de trazas más detallado para la biblioteca cliente de Kafka subyacente, librdkafka, mediante la configuración.
Manejo de mensajes malformados
Kafka se utiliza a menudo como un “vertedero” de datos. Esto hace que los topics contengan formatos de mensaje mixtos y nombres de campo incoherentes. Evite esta situación y aproveche funcionalidades de Kafka como Kafka Streams o ksqlDB para garantizar que los mensajes estén bien formados y sean coherentes antes de insertarlos en Kafka. Si estas opciones no son viables, ClickHouse ofrece algunas funciones que pueden ayudar.
  • Trate los campos del mensaje como cadenas. Se pueden usar funciones en la sentencia de la vista materializada para realizar la limpieza y la conversión de tipos si es necesario. Esto no debería considerarse una solución para producción, pero puede ser útil para una ingestión puntual.
  • Si está consumiendo JSON desde un topic con el format JSONEachRow, use la setting input_format_skip_unknown_fields. Al escribir datos, ClickHouse, de forma predeterminada, throws an exception si los datos de entrada contienen columnas que no existen en la tabla de destino. Sin embargo, si esta opción está enabled, esas columnas sobrantes se ignorarán. De nuevo, esto no es una solución apta para producción y podría confundir a otros usuarios.
  • Considere la setting kafka_skip_broken_messages. Esto requiere que el usuario especifique el nivel de tolerancia por bloque para mensajes malformados, en el contexto de kafka_max_block_size. Si se supera esta tolerancia (medida en número absoluto de mensajes), se restablecerá el comportamiento habitual de excepción y se omitirán otros mensajes.
Semántica de entrega y problemas con los duplicados
El motor de tabla Kafka tiene semántica de al menos una vez. Es posible que haya duplicados en varias circunstancias poco frecuentes ya conocidas. Por ejemplo, los mensajes pueden leerse de Kafka e insertarse correctamente en ClickHouse. Antes de que pueda confirmarse el nuevo offset, se pierde la conexión con Kafka. En esta situación, es necesario reintentar el bloque. El bloque puede deduplicarse usando una tabla distribuida o ReplicatedMergeTree como tabla de destino. Aunque esto reduce la probabilidad de filas duplicadas, depende de que los bloques sean idénticos. Eventos como un reequilibrio de Kafka pueden invalidar este supuesto y provocar duplicados en casos poco frecuentes.
Inserciones basadas en quorum
Puede que necesite inserciones basadas en quorum en los casos en los que se requieran mayores garantías de entrega en ClickHouse. Esto no se puede configurar en la vista materializada ni en la tabla de destino. Sin embargo, sí se puede configurar para perfiles de usuario, por ejemplo.

De ClickHouse a Kafka

Aunque es un caso de uso menos habitual, los datos de ClickHouse también pueden persistirse en Kafka. Por ejemplo, insertaremos filas manualmente en un motor de tabla Kafka. Estos datos los leerá el mismo motor Kafka, cuya vista materializada colocará los datos en una tabla MergeTree. Por último, demostramos cómo aplicar vistas materializadas a las inserciones en Kafka para leer tablas a partir de tablas de origen existentes.

Pasos

Nuestro objetivo inicial se ilustra mejor de la siguiente manera: Suponemos que ya ha creado las tablas y vistas en los pasos de Kafka a ClickHouse y que el topic ya se ha consumido por completo.
1. Insertar filas directamente
Primero, confirme el recuento de la tabla de destino.
Deberías tener 200.000 filas:
Ahora inserte filas de la tabla de destino de GitHub de nuevo en el motor de tabla de Kafka github_queue. Observe cómo utilizamos el formato JSONEachRow y LIMITAMOS la consulta a 100.
Vuelva a contar las filas en GitHub para confirmar que han aumentado en 100. Como se muestra en el diagrama anterior, las filas se han insertado en Kafka mediante el motor de tabla Kafka antes de volver a leerse con el mismo motor e insertarse en la tabla de destino GitHub mediante nuestra vista materializada.
Deberías ver 100 filas más:
2. Uso de vistas materializadas
Podemos utilizar vistas materializadas para enviar mensajes a un motor Kafka (y a un topic) cuando se insertan documentos en una tabla. Cuando se insertan filas en la tabla GitHub, se activa una vista materializada, lo que hace que las filas se vuelvan a insertar en un motor Kafka y en un topic nuevo. De nuevo, se ilustra mejor de la siguiente manera: Cree un topic de Kafka nuevo github_out o uno equivalente. Asegúrese de que un motor de tabla Kafka github_out_queue apunte a este topic.
Ahora cree una nueva vista materializada github_out_mv para que apunte a la tabla de GitHub e inserte filas en el engine anterior cuando se active. Como resultado, las nuevas filas añadidas a la tabla de GitHub se enviarán a nuestro nuevo topic de Kafka.
Si insertas datos en el topic original de github, creado como parte de Kafka to ClickHouse, los documentos aparecerán mágicamente en el topic “github_clickhouse”. Confírmalo con las herramientas nativas de Kafka. Por ejemplo, a continuación, insertamos 100 filas en el topic github con kcat para un topic alojado en Confluent Cloud:
Una lectura en el topic github_out debería confirmar la entrega de los mensajes.
Aunque se trata de un ejemplo complejo, ilustra la potencia de las vistas materializadas cuando se utilizan junto con el motor Kafka.

Clústeres y rendimiento

Trabajo con clústeres de ClickHouse

Mediante grupos de consumidores de Kafka, varias instancias de ClickHouse pueden leer potencialmente del mismo topic. A cada consumidor se le asigna una partición del topic en una relación 1:1. Al escalar el consumo de ClickHouse con el motor de tabla Kafka, tenga en cuenta que el número total de consumidores dentro de un clúster no puede superar el número de particiones del topic. Por lo tanto, asegúrese de configurar de antemano un particionado adecuado para el topic. Se pueden configurar varias instancias de ClickHouse para leer de un topic con el mismo ID de grupo de consumidores, especificado durante la creación del motor de tabla Kafka. De este modo, cada instancia leerá de una o más particiones e insertará segmentos en su tabla de destino local. A su vez, las tablas de destino pueden configurarse para usar ReplicatedMergeTree y gestionar la duplicación de datos. Este enfoque permite escalar las lecturas de Kafka junto con el clúster de ClickHouse, siempre que haya suficientes particiones de Kafka.

Ajuste del rendimiento

Tenga en cuenta lo siguiente si busca aumentar el rendimiento de una tabla con el motor Kafka:
  • El rendimiento variará según el tamaño de los mensajes, el formato y los tipos de las tablas de destino. Un rendimiento de 100k filas/s en un único motor de tabla debe considerarse alcanzable. De forma predeterminada, los mensajes se leen en bloques, controlados por el parámetro kafka_max_block_size. Por defecto, este se establece en max_insert_block_size, cuyo valor predeterminado es 1,048,576. Salvo que los mensajes sean extremadamente grandes, casi siempre conviene aumentarlo. No es raro usar valores entre 500k y 1M. Pruebe y evalúe el efecto sobre el rendimiento.
  • El número de consumidores de un motor de tabla puede aumentarse con kafka_num_consumers. Sin embargo, de forma predeterminada, las inserciones se serializan en un único hilo, a menos que kafka_thread_per_consumer se cambie desde su valor predeterminado de 1. Establézcalo en 1 para garantizar que los flush se realicen en paralelo. Tenga en cuenta que crear una tabla con el motor Kafka con N consumidores (y kafka_thread_per_consumer=1) es lógicamente equivalente a crear N motores Kafka, cada uno con una vista materializada y kafka_thread_per_consumer=0.
  • Aumentar el número de consumidores no sale gratis. Cada consumidor mantiene sus propios búferes e hilos, lo que incrementa la sobrecarga del servidor. Tenga presente esta sobrecarga y, si es posible, escale primero de forma lineal en todo el clúster.
  • Si el rendimiento de los mensajes de Kafka es variable y los retrasos son aceptables, considere aumentar stream_flush_interval_ms para asegurarse de que se escriban bloques más grandes.
  • background_message_broker_schedule_pool_size establece el número de hilos que realizan tareas en segundo plano. Estos hilos se usan para el streaming de Kafka. Esta configuración se aplica al iniciar el servidor ClickHouse y no puede cambiarse en una sesión de usuario; su valor predeterminado es 16. Si ve timeouts en los logs, puede ser conveniente aumentar este valor.
  • Para la comunicación con Kafka se utiliza la biblioteca librdkafka, que a su vez crea hilos. Por tanto, un gran número de tablas Kafka o de consumidores puede dar lugar a una gran cantidad de cambios de contexto. Distribuya esta carga en todo el clúster, replicando solo las tablas de destino si es posible, o considere usar un motor de tabla para leer de varios topics; se admite una lista de valores. Varias vistas materializadas pueden leer de una sola tabla, cada una filtrando los datos de un topic concreto.
Cualquier cambio en la configuración debe probarse. Recomendamos supervisar el consumer lag de Kafka para asegurarse de que el escalado sea el adecuado.

Ajustes adicionales

Además de los ajustes mencionados anteriormente, los siguientes pueden resultar de interés:
  • Kafka_max_wait_ms - El tiempo de espera, en milisegundos, para leer mensajes de Kafka antes de reintentar. Se establece a nivel de perfil de usuario y su valor predeterminado es 5000.
Todos los ajustes de la librdkafka subyacente también pueden colocarse en los archivos de configuración de ClickHouse dentro de un elemento kafka; los nombres de los ajustes deben ser elementos XML con los puntos sustituidos por guiones bajos; p. ej.
Estas son configuraciones avanzadas y le sugerimos consultar la documentación de Kafka para obtener una explicación más detallada.
Última modificación el 19 de junio de 2026