Nuestro setup de ClickHouse para escalar
ObsessionDB escala ClickHouse® en los dos lados de la carga. En ingesta, una migración cargó cientos de terabytes en pocos días a más de 300 inserts por segundo, y un solo nodo de 32 CPU sostuvo 31,8 millones de filas por segundo durante 90 minutos. En servicio, una tabla de 203.000 millones de filas y 56 TB responde lookups aleatorios por debajo del segundo en p99. Debajo hay una sola copia de los datos en object storage, compute sin estado, merges con localidad y settings que cualquier cluster de ClickHouse puede usar.
La semana pasada publicamos nuestro setup para reducir la latencia. Este artículo es el otro eje del escalado de ClickHouse: qué le pasa al mismo cluster cuando los datos crecen diez veces y el ritmo de ingesta crece con ellos.
Las dos curvas salen de una serie de decisiones sobre dónde viven los bytes y qué nodo hace cada trabajo. La primera mitad del artículo recorre esa pila capa por capa, cada una comparada con lo que hace un cluster autogestionado en la misma situación. La segunda mitad es el kit: los settings que mueven la escala en cualquier cluster, separados en ingesta, merges y query.
El setup, capa por capa
Añadir un nodo a uno de nuestros clusters no mueve datos. Casi todas las capas de abajo descansan sobre esa única propiedad, y es la que un cluster autogestionado no puede comprar con settings.
Una sola copia de los datos, compute sin estado
Alloy es nuestro motor de almacenamiento, desarrollado contra la API de SharedMergeTree. Los metadatos de una tabla viven en una capa de coordinación, sus datos viven en object storage como objetos inmutables, y los nodos de compute no guardan nada duradero. Un nodo que entra se suscribe a los metadatos y empieza a servir; su parte del anillo de cache se llena con el primer acceso. Un nodo que sale cede sus claves al siguiente nodo del ranking y el anillo se recompone. Nuestra documentación de escalado lo resume en una línea: cambiar la forma del cluster no mueve datos, sin re-sharding y sin rebalanceo. Construir sobre un ClickHouse desacoplado cuenta qué cambia para las tablas de encima.
ReplicatedMergeTree, el motor que corre un cluster autogestionado, guarda una copia completa de los bytes en cada réplica y las coordina a través de Keeper. Tres réplicas de una tabla de 54 GiB ocupan unos 162 GiB. Cada tabla lógica se convierte en una tabla _local por shard más una fachada Distributed, y cada ALTER sale ON CLUSTER, donde puede aplicarse en unos shards y fallar en otros. Nosotros funcionamos así durante un año; añadir un shard significaba copiar más de 25 TB.
Poner un disco S3 debajo de MergeTree no cambia eso. Con una storage policy por niveles, cada réplica sigue guardando su propia copia de los ficheros en S3, y los metadatos de las parts se quedan en el disco local de cada nodo. La replicación zero-copy comparte los objetos pero mantiene los metadatos en local, y la documentación de ClickHouse la marca como no lista para producción. Viene desactivada por defecto desde la 22.8. El almacenamiento por niveles es un nivel frío con el mismo acoplamiento.
Compute separado sobre los mismos objetos
Un cluster solía significar una sola mezcla de cargas. Sin nada duradero en el nodo, un segundo grupo de nodos puede leer los mismos objetos, y los dos grupos no compiten por memoria ni por ancho de banda. Los nodos que sirven a un cliente van en un grupo y, si hace falta, un nodo aparte para sus backfills y reconstrucciones pesadas va en otro. Los trabajos de reingesta funcionan igual, en infraestructura aislada en vez de en los nodos que sirven la API.
La versión extrema es un datashare: el compute de otra organización lee tus objetos, en solo lectura, sin copia ni réplica, y ninguna de las dos partes puede quemar los recursos de la otra.
La cache distribuida
Todas las lecturas pasan por la cache distribuida. Cada nodo aporta su NVMe a un único anillo con rendezvous hashing, y una part se cachea en el momento en que se escribe, así que las particiones más recientes se sirven calientes desde su primera lectura. Un salto entre nodos está por debajo del milisegundo, mientras que object storage se mueve en decenas de milisegundos en la mediana y cientos en la cola, como midió el artículo de latencia. Su ancho de banda crece con el cluster: 20 nodos dan 20 veces la red de un nodo, donde una cache centralizada se queda cerca de tres. Dimensiónala para que el working set por nodo quepa en un 80% del disco de cache de ese nodo (la documentación da la fórmula), o tu latencia seguirá a object storage. Ese artículo también cubre cómo se comporta la cache bajo carga.
Merges que se quedan donde están los datos
El 17% es la proporción de merges que un cluster de seis nodos sobre storage compartido haría en local si las parts cayeran al azar: una posibilidad entre seis de que las entradas estén en el nodo que hace el trabajo. Nuestros clusters más ocupados rondan el 70% en local. Calculamos la diferencia en unas cuatro veces menos ancho de banda de merge con seis nodos, y el benchmark de 10.000 millones de filas es donde se ve ese margen: hacemos merges de parts de hasta 150 GB sin congestionar la red.
El mecanismo no tiene scheduler. En un cluster sin estado, cualquier nodo puede hacer merge de cualquier conjunto de parts, porque las parts viven en object storage. Cada nodo ejecuta el mismo selector de merges y aplica el mismo ranking de rendezvous hashing a una única pregunta local: si es el nodo mejor situado para ese conjunto de candidatas. El nodo que sale más alto en el ranking para las entradas se queda el merge, y su salida aterriza en su propietario según el ranking, que es adonde se enrutan las siguientes lecturas de esa part. La documentación describe la colocación como algo que emerge de que todos los nodos apliquen el mismo ranking a los mismos hechos. Hay dos settings, y no hace falta tocarlos: cache_locality_aware_merges está activado, y min_bytes_for_locality_aware_merges vale 100 MB por defecto para que las parts pequeñas, que son las más numerosas, se fusionen donde haya un hueco libre en vez de amontonarse en un nodo.
Un merge suele leer varios cientos de parts de entrada. En un cluster con colocación aleatoria, la proporción de esas entradas en la cache del nodo que hace el merge es 1/N, así que cada nodo que añades mueve más bytes por la red para hacer el mismo trabajo de merge. La localidad mantiene la proporción local cerca del 70% a medida que N crece. Esa es toda la razón de que añadir nodos ayude a los merges aquí y los perjudique en un cluster sobre storage compartido que coloca los merges al azar.
Queries que se reparten hacia donde los datos están calientes
Nuestro fan-out de parallel replicas se diferencia del ClickHouse estándar en el enrutado. El planificador de lecturas reutiliza el ranking de rendezvous hashing que coloca los merges: cada sub-query toma primero los rangos que posee su propio nodo en el anillo de cache, y después roba los rangos que quedan al siguiente nodo del ranking en vez de a uno al azar, así que la tasa de aciertos se mantiene aunque el propietario esté ocupado. En la tabla de 56 TB, en un caso idealizado de lectura por red, la misma query tardó 1,3 s en un nodo, 1,0 s en seis nodos con parallel replicas a secas, y 0,5 s en esos seis nodos con enrutado por localidad.
El ClickHouse open source 26.6 añadió un segundo modo de ejecución experimental, la ejecución distribuida multi-stage, y la 26.8 le añadió un optimizador basado en coste. El planificador parte una query en etapas unidas por intercambios scatter, broadcast, gather y shuffle, así que un GROUP BY de alta cardinalidad ya no hace pasar todos los resultados parciales por la memoria de un único coordinador. Lo estamos probando en la 26.8. En nuestro cluster de pruebas de tres nodos, en septiembre, una agregación de alta cardinalidad corrió 2,86 veces más rápido que en un nodo, y otra más pesada que falla en un solo nodo y bajo parallel replicas termina con ejecución multi-stage y con bastante menos memoria.
El kit: qué mueve la escala en cualquier cluster de ClickHouse
Las cinco capas vienen con la plataforma. Los settings de abajo son los que están en tu mano, en cualquier cluster, y deciden si las capas llegan a ayudar.
Ingesta
Agrupa los inserts. ClickHouse recomienda en torno a un insert por segundo y tabla, con miles de filas en cada uno. Muchos inserts pequeños hacen muchas parts pequeñas, y cada part hay que fusionarla después. Los async inserts hacen ese agrupado en el servidor. Mantén wait_for_async_insert = 1 para que un insert fallido falle a la vista, y dale al cliente un timeout más largo del que necesita el servidor, o un reintento escribirá las mismas filas dos veces. Desde la 26.8, max_insert_threads usa todos los cores por defecto, lo que produce más parts y quita CPU a los merges; fíjalo bajo para las cargas masivas.
| Setting | Qué te da | Qué te cuesta |
|---|---|---|
| Agrupar, ~1 insert/s por tabla | Menos parts, menos merges | Algo de latencia en el cliente |
async_insert = 1, wait_for_async_insert = 1 | Agrupado en el servidor | Un timeout que tienes que dimensionar |
max_insert_threads (auto desde la 26.8) | Cargas masivas más rápidas | Más parts, menos CPU para merges |
deduplicate_insert = enable | Los reintentos no duplican filas | Solo dentro de la ventana de dedup |
Merges
Un merge reescribe un conjunto de parts en una sola, y sigue hasta que las parts llegan a unos 150 GB. Un nodo estándar ejecuta hasta 32 merges a la vez y elige primero las pequeñas, así que un número bajo de parts hace más por ti que cualquier setting del pool. Dos cosas que aprendimos por las malas: OPTIMIZE FINAL sobre parts que pasan de 150 GB puede fallar, y FINAL se ralentiza con cada part que tiene que leer, así que do_not_merge_across_partitions_select_final = 1 es una victoria barata en tablas particionadas.
| Setting | Qué te da | Qué te cuesta |
|---|---|---|
background_pool_size | Más merges en vuelo | CPU y memoria que se quitan a las queries |
parts_to_delay_insert / parts_to_throw_insert | Un aviso antes de Too many parts | Subirlos solo aplaza el trabajo |
min_age_to_force_merge_seconds | Las particiones tranquilas también compactan | Reescrituras programadas |
do_not_merge_across_partitions_select_final | FINAL más barato | Nada en una tabla particionada |
Query
Arregla el patrón de acceso antes de añadir un nodo. Un cliente nos pidió dos nodos más; una query escaneaba un petabyte al día, y un solo índice recortó la CPU del cluster un 99,5%. max_threads usa todos los cores por defecto; en lookups bajo inserts pesados, fijarlo a 4 u 8 por query nos dio un p99 estable. Las parallel replicas ayudan a los scans grandes y perjudican a los point lookups, así que actívalas por query con parallel_replicas_min_number_of_rows_per_replica como barandilla. Después de añadir nodos, lanza SYSTEM PREWARM para que los datos estén calientes en la cache.
| Setting | Qué te da | Qué te cuesta |
|---|---|---|
| Una sorting key que case con el filtro | Salta gránulos en vez de escanear | Un rediseño, una vez |
max_threads de 4 a 8 por query de lookup | p99 estable bajo carga mixta | Velocidad punta en scans anchos |
enable_parallel_replicas por query | Los scans grandes usan todo el cluster | Sobrecoste en queries pequeñas |
SYSTEM PREWARM tras un redimensionado | Los nodos nuevos arrancan calientes | Tiempo de red por adelantado |
Qué suman las capas
Las capas se notan en los tres ejes por los que se juzga un cluster de ClickHouse: lo rápido que traga datos, lo rápido que los procesa y lo rápido que responde.
Primero la ingesta, porque es donde un cluster replicado se queda sin carretera antes. En ReplicatedMergeTree un insert se escribe una vez y luego se copia a cada réplica, unas diez entradas en Keeper por INSERT, con el cluster entero presupuestado en unos pocos cientos de inserts por segundo, y cada réplica repite cada merge sobre su propia copia. Aquí un insert es una sola escritura: el nodo que lo recibe escribe la part una vez, a través de la cache distribuida hacia object storage, entrega los metadatos a la capa de coordinación, y el resto de nodos reciben un aviso de part nueva y precalientan el índice sin copiar un byte. El packed storage mete las columnas de una part en un único objeto, unas 15 veces menos escrituras en parts anchas. Los números vienen detrás. Un nodo de 32 CPU sostuvo 31,8 millones de filas por segundo durante 90 minutos este junio, cerca de un millón de filas por segundo y CPU, con los merges corriendo todo el tiempo. La migración detrás de el artículo de Numia cargó cientos de terabytes de más de veinte cadenas en pocos días a más de 300 inserts por segundo y tabla, y los clientes que alimentaban el cluster fueron lo que lo limitó; una hora de pico en el mismo cluster en mayo produjo 335.000 parts nuevas y entre dos y tres millones de peticiones a object storage sin un solo error de merge ni de storage. El backfill de un cliente más avanzado el verano corrió a cinco o seis millones de filas por segundo en nueve nodos, y dos clientes en prueba de concepto, los dos con experiencia en ClickHouse Cloud, nos dijeron que no habían visto antes ese ritmo de ingesta.
Los merges son donde añadir nodos paga o castiga. Como cualquier nodo puede hacer merge de cualquier conjunto de parts, la capacidad de merge crece con el cluster en vez de repetirse en cada réplica, y el enrutado por localidad mantiene cerca del 70% de las entradas de un merge en el nodo que hace el trabajo, donde la colocación aleatoria en seis nodos daría un 17%: unas cuatro veces menos ancho de banda de merge. Ese margen es lo que nos deja hacer merges de parts de hasta 150 GB sin congestionar la red, y es la razón de que los ritmos de ingesta de arriba se sostengan mientras corre la compactación que hay detrás.
Las queries reciben la misma copia única de los datos, caliente desde el momento en que se escribe, con las lecturas enrutadas al nodo que la tiene. La tabla de 56 TB responde lookups por dirección aleatoria con un p99 de 703 ms en seis nodos, y el dataset de 18 TB servido a través de materialized views quedó por debajo de 410 ms en p99 en 25.000 queries con cinco. Las parallel replicas se llevan los scans pesados, 16,3 veces en la query más pesada del benchmark de 10.000 millones de filas, solo por query, y la ejecución multi-stage en la 26.8 se lleva las agregaciones de alta cardinalidad que ni un solo nodo ni las parallel replicas manejan bien.
El último eje es lo que te cuesta cambiar de forma. Añadir o quitar un nodo no mueve datos, un trabajo pesado puede correr en su propio grupo de compute sin que la API se entere, y la factura es compute más storage sin cargo por query, que es por lo que la curva de coste por terabyte de la intro baja a medida que crecen los datos.
Si algo no cuadra, recorre esta tabla hacia abajo. Cada fila es más barata de comprobar que la siguiente, y la mayoría de los tickets se cierran antes de llegar al final.
| Síntoma | Comprueba primero | El movimiento |
|---|---|---|
| Inserts frenados o rechazados | Parts por partición en system.parts | Agrupa; después sube los umbrales en esa tabla y paga la deuda de merges más tarde |
| Los merges se retrasan más a medida que añades nodos | Volumen de fetches entre nodos durante las ventanas de merge | Merges con localidad, o menos nodos y más grandes |
| Los scans grandes ignoran tu número de nodos | enable_parallel_replicas por query, analyzer activado | Reparte solo las formas pesadas, mantén la barandilla de filas mínimas; en la 26.8, ejecución multi-stage para agregaciones de alta cardinalidad |
| Un backfill ralentiza la API | En qué nodos corre el trabajo | Compute separado sobre los mismos objetos |
| Los reintentos cuentan doble | La ventana de dedup frente a tu intervalo de reintento | Dedup por bloque en los inserts; vistas idempotentes para el fan-out |
| Todo en verde y aun así lento con N nodos | Peticiones a S3 por query, working set de la cache frente al disco | Estás en la línea de la arquitectura; ningún setting la mueve |
La última pregunta es la forma del propio cluster. Nuestra regla de dimensionado: las operaciones pesadas de tipo warehouse quieren menos nodos con recursos más concentrados, porque un merge grande o un join grande necesitan la memoria de un solo nodo; el servicio con SLA ajustado y demanda elástica quiere una flota mayor, por el disco de cache y la capacidad de scan en paralelo; en caso de duda, menos nodos y más grandes, y no menos de tres. Un cluster de dos nodos pierde la mitad de su capacidad en un rolling update. Uno de nueve pierde un 11%.
La promesa a la que sirve todo esto es la misma en los dos ejes: menos hardware para el mismo rendimiento, o más rendimiento con el mismo hardware. Si tu cluster no se comporta así, lo miramos contigo. Nuestra performance audit coge tu query log y tus recuentos de parts, les pasa esta checklist y las más profundas, y te devuelve lo que encuentre, vayas a correr en ObsessionDB o no.
FAQ
En ReplicatedMergeTree una réplica nueva repite cada merge; añade concurrencia y acelera una query solo si esta se reparte. Sobre storage compartido sin enrutado por localidad, la proporción de entradas de un merge que ya están en local cae como 1/N, así que cada nodo que añades mueve más bytes de merge. Revisa primero el patrón de acceso; nuestra última petición de "añadid nodos" era un índice que faltaba.
La documentación recomienda en torno a un insert por segundo y tabla, en lotes de 10.000 a 100.000 filas. Cada insert crea al menos una part que hay que fusionar, así que el techo es la capacidad de merge. Hemos llevado clientes por encima de 300 inserts por segundo y tabla, pero los lotes eran grandes y el pool de merges tenía margen.
Los inserts crean parts más rápido de lo que los merges en segundo plano las combinan, y cuando una partición pasa de parts_to_throw_insert (3.000 en open source, 600 en el nuestro) los inserts se rechazan. La causa suele ser inserts pequeños y frecuentes o demasiadas particiones. Agrupa primero, revisa la clave de partición después, y sube el umbral solo como puente mientras los merges se ponen al día.
Para scans pesados y agregaciones grandes, donde la lectura se puede repartir entre nodos; el benchmark de 10.000 millones de filas llegó a 16,3 veces con 20 nodos. No en global: un point lookup pasó de 0,5 s a 37 s en la misma prueba. Actívalas por query, mantén parallel_replicas_min_number_of_rows_per_replica como barandilla, y recuerda que las projections se saltan bajo lecturas en paralelo.
No. Una storage policy por niveles mueve las parts frías a un disco S3, pero cada réplica sigue guardando su propia copia de los objetos y sus propios metadatos de parts en disco local, y las réplicas siguen coordinándose a través de Keeper. Separar significa una sola copia de los datos, los metadatos fuera de los nodos, y compute que se puede añadir o sustituir sin mover nada.
Sí, en tablas replicadas. Cada bloque insertado recibe un hash que se guarda en Keeper, y un bloque reintentado con el mismo hash se descarta; open source conserva los últimos 10.000 hashes durante una hora. Los destinos de las materialized views solo quedan cubiertos cuando deduplicate_blocks_in_dependent_materialized_views está activado, y los async inserts a través de una vista que emite varios bloques lanzan NOT_IMPLEMENTED.
Seguir Leyendo
Publicado originalmente en obsessionDB. Lee el artículo original aquí.
ClickHouse is a registered trademark of ClickHouse, Inc. https://clickhouse.com