Introduction
Dans ce guide, nous verrons d’abord comment ClickHouse répartit une requête sur plusieurs shards via des tables distribuées, puis comment une requête peut tirer parti de plusieurs répliques pour son exécution.Dans une architecture shared-nothing, les clusters sont généralement répartis en plusieurs shards, chaque shard contenant un sous-ensemble de l’ensemble des données. Une table distribuée se place au-dessus de ces shards et fournit une vue unifiée de l’intégralité des données. Les lectures peuvent être envoyées à la table locale. L’exécution de la requête n’aura lieu que sur le shard spécifié, ou bien elle peut être envoyée à la table distribuée, auquel cas chaque shard exécutera la requête donnée. Le serveur sur lequel la table distribuée a été interrogée agrégera les données et renverra la réponse au client : La figure ci-dessus illustre ce qui se passe lorsqu’un client interroge une table distribuée :
- La requête SELECT est envoyée à une table distribuée sur un nœud choisi arbitrairement (via une stratégie round-robin ou après avoir été redirigée vers un serveur spécifique par un répartiteur de charge). Ce nœud va alors jouer le rôle de coordinateur.
- Le nœud localise chaque shard qui doit exécuter la requête à l’aide des informations spécifiées par la table distribuée, puis la requête est envoyée à chaque shard.
- Chaque shard lit, filtre et agrège les données localement, puis renvoie un état fusionnable au coordinateur.
- Le nœud coordinateur fusionne les données, puis renvoie la réponse au client.
Architecture non shardée
ClickHouse Cloud présente une architecture très différente de celle décrite ci-dessus. (Voir “Architecture de ClickHouse Cloud” pour plus de détails). Avec la séparation du calcul et du stockage, ainsi qu’une capacité de stockage virtuellement infinie, le besoin de shards devient moins important. La figure ci-dessous montre l’architecture de ClickHouse Cloud : Cette architecture permet d’ajouter et de supprimer des répliques presque instantanément, ce qui garantit une très forte scalabilité du cluster. Le cluster ClickHouse Keeper (illustré à droite) garantit l’existence d’une source de vérité unique pour les métadonnées. Les répliques peuvent récupérer les métadonnées depuis le cluster ClickHouse Keeper et conservent toutes les mêmes données. Les données elles-mêmes sont stockées dans le stockage objet, et le cache SSD permet d’accélérer les requêtes. Mais comment distribuer désormais l’exécution des requêtes sur plusieurs serveurs ? Dans une architecture shardée, c’était assez évident puisque chaque shard pouvait effectivement exécuter une requête sur un sous-ensemble des données. Comment cela fonctionne-t-il lorsqu’il n’y a pas de sharding ?Présentation des répliques parallèles
Pour paralléliser l’exécution des requêtes sur plusieurs serveurs, nous devons d’abord être en mesure de désigner l’un de nos serveurs comme coordinateur. Le coordinateur est celui qui crée la liste des tâches à exécuter, veille à ce qu’elles soient toutes exécutées et agrégées, puis à ce que le résultat soit renvoyé au client. Comme dans la plupart des systèmes distribués, ce sera le rôle du nœud qui reçoit la requête initiale. Nous devons également définir l’unité de travail. Dans une architecture shardée, l’unité de travail est le shard, c’est-à-dire un sous-ensemble des données. Avec les répliques parallèles, nous utiliserons une petite portion de la table, appelée granules, comme unité de travail. Voyons maintenant comment cela fonctionne en pratique à l’aide de la figure ci-dessous : Avec les répliques parallèles :- La requête du client est envoyée à un nœud après être passée par un répartiteur de charge. Ce nœud devient le coordinateur de cette requête.
- Le nœud analyse l’index de chaque part et sélectionne les bonnes parts et granules à traiter.
- Le coordinateur répartit la charge de travail en un ensemble de granules pouvant être attribués à différentes répliques.
- Chaque ensemble de granules est traité par les répliques correspondantes, puis un état fusionnable est envoyé au coordinateur une fois le traitement terminé.
- Enfin, le coordinateur fusionne tous les résultats des répliques puis renvoie la réponse au client.
- Certaines répliques peuvent être indisponibles.
- La réplication dans ClickHouse est asynchrone ; à un instant donné, certaines répliques peuvent ne pas avoir les mêmes parts.
- Il faut également gérer la tail latency.
- Le cache du système de fichiers varie d’une réplique à l’autre en fonction de l’activité sur chacune d’elles, ce qui signifie qu’une attribution aléatoire des tâches pourrait entraîner des performances moins bonnes en raison de la localité du cache.
Annonces
Pour répondre aux points (1) et (2) de la liste ci-dessus, nous avons introduit le concept d’annonce. Visualisons son fonctionnement à l’aide de la figure ci-dessous :- La requête du client est envoyée à un nœud après être passée par un répartiteur de charge. Ce nœud devient le coordinateur de cette requête.
- Le nœud coordinateur envoie ensuite une requête pour récupérer les annonces de toutes les répliques du cluster. Les répliques peuvent avoir des vues légèrement différentes de l’ensemble actuel des parts d’une table. Il faut donc collecter ces informations afin d’éviter des décisions de planification incorrectes.
- Le nœud coordinateur utilise ensuite les annonces pour définir un ensemble de granules pouvant être attribuées aux différentes répliques. Ici, par exemple, on peut voir qu’aucune granule de la part 3 n’a été attribuée à la réplique 2, car cette réplique n’a pas indiqué cette part dans son annonce. Notez également qu’aucune tâche n’a été attribuée à la réplique 3, car la réplique n’a fourni aucune annonce.
- Une fois que chaque réplique a traité la requête sur son sous-ensemble de granules et que l’état fusionnable a été renvoyé au coordinateur, celui-ci fusionne les résultats et renvoie la réponse au client.
Coordination dynamique
Pour remédier au problème de latence des requêtes les plus lentes, nous avons ajouté la coordination dynamique. Cela signifie que toutes les granules ne sont pas envoyées à une réplique en une seule requête, mais que chaque réplique peut demander une nouvelle tâche (un ensemble de granules à traiter) au coordinateur. Le coordinateur fournit alors à la réplique l’ensemble de granules en fonction de l’annonce reçue. Supposons que nous soyons à l’étape du processus où toutes les répliques ont envoyé une annonce avec toutes les parts. La figure ci-dessous illustre le fonctionnement de la coordination dynamique :- Les répliques indiquent au nœud coordinateur qu’elles sont en mesure de traiter des tâches, et peuvent également préciser la quantité de travail qu’elles peuvent prendre en charge.
- Le coordinateur attribue des tâches aux répliques.
- Les répliques 1 et 2 terminent leur tâche très rapidement. Elles demandent alors une autre tâche au nœud coordinateur.
- Le coordinateur attribue de nouvelles tâches aux répliques 1 et 2.
- Toutes les répliques ont maintenant terminé le traitement de leur tâche. Elles demandent davantage de tâches.
- Le coordinateur, à l’aide des annonces, vérifie quelles tâches restent à traiter, mais il n’en reste aucune.
- Le coordinateur indique aux répliques que tout a été traité. Il va maintenant fusionner tous les états fusionnables et répondre à la requête.
Gestion de la localité du cache
| Réplique 1 | Réplique 2 | Réplique 3 | |
|---|---|---|---|
| Part 1 | g1, g6, g7 | g2, g4, g5 | g3 |
| Part 2 | g1 | g2, g4, g5 | g3 |
| Part 3 | g1, g6 | g2, g4, g5 | g3 |
max_parallel_replicas est inférieur
au nombre de répliques, des répliques aléatoires sont alors sélectionnées pour l’exécution des requêtes.
Vol de tâches
Limitations
Cette fonctionnalité présente des limitations connues, dont les principales sont documentées dans cette section.Si vous trouvez un ticket qui ne figure pas parmi les limitations indiquées ci-dessous et
soupçonnez les répliques parallèles d’en être la cause, veuillez ouvrir un ticket sur GitHub avec
le label
comp-parallel-replicas.Investiguer les problèmes liés aux répliques parallèles
Vous pouvez vérifier quels paramètres sont utilisés pour chaque requête dans la tablesystem.query_log. Vous pouvez
également consulter la table system.events
pour voir tous les événements survenus sur le serveur, et vous pouvez utiliser la
table function clusterAllReplicas pour voir les tables sur tous les répliques
(si vous utilisez ClickHouse Cloud, utilisez default).
Query
Réponse
Réponse
Response
system.text_log contient également des informations sur l’exécution de requêtes utilisant des répliques parallèles :
Query
Réponse
Réponse
Response
EXPLAIN PIPELINE. Il met en évidence la manière dont ClickHouse
va exécuter une requête et quelles ressources seront utilisées pour
son exécution. Prenons la requête suivante comme exemple :
EXPLAIN PIPELINE (without parallel replica)
EXPLAIN PIPELINE (with parallel replica)