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.
Resumen
Pasos
1
Prepare
Si tiene datos cargados en un topic de destino, puede adaptar lo siguiente para usarlo en su conjunto de datos. Como alternativa, se proporciona un conjunto de datos de muestra de GitHub aquí. Este conjunto de datos se usa en los ejemplos siguientes y utiliza 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í, por brevedad. Aun así, sigue siendo suficiente para que funcionen la mayoría de las consultas publicadas con el conjunto de datos.
2
Configure ClickHouse
Este paso es obligatorio si se está conectando a un Kafka seguro. Estos ajustes no pueden pasarse mediante comandos SQL DDL y deben configurarse en el archivo config.xml de ClickHouse. Suponemos que se está conectando a una instancia protegida con SASL. Este es el método más sencillo para interactuar con Confluent Cloud.Coloque el fragmento anterior en un archivo nuevo dentro de su directorio conf.d/ o combínelo con los archivos de configuración existentes. Para conocer los ajustes que se pueden configurar, consulte aquí.También vamos a crear una base de datos llamada Una vez creada la base de datos, deberá cambiar a ella:
KafkaEngine para usar en este tutorial: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
Cree el topic y llénelo
A continuación, crearemos un topic. Hay varias herramientas que podemos usar para hacerlo. Si ejecutamos Kafka localmente en nuestra máquina o dentro de un contenedor Docker, RPK funciona bien. Podemos crear un topic llamado Si ejecutamos Kafka en Confluent Cloud, quizá prefiramos usar la CLI de Confluent:Ahora necesitamos poblar 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 la ingesta debería completarse en apenas unos segundos. Si desea trabajar con un conjunto de datos más grande, consulte la sección de conjuntos de datos grandes del repositorio de GitHub ClickHouse/kafka-samples.
github con 5 particiones ejecutando el siguiente comando:5
Cree el motor de tabla de Kafka
El siguiente ejemplo crea un motor de tabla con el mismo esquema que la tabla MergeTree. Esto no es estrictamente necesario, ya que puede tener un alias o columnas efímeras en la tabla de destino. Sin embargo, la configuración es importante; observe el uso de A continuación, analizamos la configuración del motor y el ajuste del rendimiento. En este punto, un
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 del grupo de consumidores, respectivamente. Los topics pueden ser, de hecho, una lista de valores.select sencillo en la tabla github_queue debería leer algunas filas. Tenga en cuenta que esto hará avanzar los offsets del consumidor, lo que impedirá volver a leer estas filas sin un reinicio. Tenga en cuenta el límite y el parámetro obligatorio stream_like_engine_allow_direct_select.6
Cree la vista materializada
La vista materializada conectará las dos tablas creadas previamente, leyendo datos del motor de tabla Kafka e insertándolos en la tabla MergeTree de destino. Podemos realizar varias transformaciones de datos. Haremos una lectura e inserción sencillas. El uso de * supone que los nombres de las columnas son idénticos (distingue entre mayúsculas y minúsculas).En el momento de su creación, la vista materializada se conecta al motor Kafka y comienza a leer, insertando filas en la tabla de destino. Este proceso continuará indefinidamente y consumirá los mensajes posteriores insertados en Kafka. No dude en volver a ejecutar el script de inserción para insertar más mensajes en Kafka.
7
Confirme que se hayan insertado filas
Compruebe que hay datos en la tabla de destino:Debería ver 200.000 filas:
Operaciones comunes
Detener & reiniciar el consumo de mensajes
Añadir metadatos de Kafka
_.
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.
Modificar la configuración del motor Kafka
Depuración de incidencias
Manejo de mensajes malformados
- 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
Inserciones basadas en quorum
De ClickHouse a Kafka
Pasos
1
Insertar filas directamente
Primero, confirme el número de filas de la tabla de destino.Debería haber 200.000 filas:Ahora inserte filas desde la tabla de destino de GitHub de nuevo en el motor de tabla Kafka github_queue. Observe cómo utilizamos el formato JSONEachRow y aplicamos LIMIT a la consulta SELECT para devolver 100 filas.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 ser leídas por el mismo motor e insertadas en la tabla de destino de GitHub por 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 tema de Kafka nuevo Ahora cree una nueva vista materializada 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:La lectura del topic Aunque es un ejemplo elaborado, ilustra el potencial de las vistas materializadas cuando se usan en combinación con el motor Kafka.
github_out o uno equivalente. Asegúrese de que un motor de tabla de Kafka github_out_queue apunte a este tema.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.github_out debería confirmar la entrega de los mensajes.Clústeres y rendimiento
Trabajo con clústeres de ClickHouse
Ajuste del rendimiento
- 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.
Ajustes adicionales
- 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.