Comprendre l'Architecture Distribuée de TiDB

Contrairement aux bases de données monolithiques, TiDB se distingue par plusieurs avantages fondamentaux :

  • Une conception entièrement distribuée offrant une scalabilité horizontale et une élasticité pour l'expansion ou la réduction des ressources.
  • La prise en charge du SQL avec une exposition du protocole réseau MySQL, assurant une compatibilité étendue avec la syntaxe MySQL pour un remplacement direct dans de nombreux scénarios.
  • Une haute disponibilité intégrée, permettant au système de réparer les données et de basculer automatiquement en cas de défaillance de quelques réplicas, de manière transparente pour les applications.
  • Le support des transactions ACID, essentiel pour les applications nécessitant une forte cohérence, comme les systèmes bancaires.
  • Un écosystème d'outils riche couvrant la migration, la synchronisation et la sauvegarde des données.

Au cœur de sa conception, la base de données distribuée TiDB est structurée en plusieurs modules interconnectés qui collaborent pour former un système complet.

  • Serveur TiDB : Il représente la couche SQL, exposant un point de terminaison de connexion compatible avec le protocole MySQL. Sa responsabilité inclut la réception des requêtes client, l'analyse syntaxique et l'optimisation des requêtes SQL, menant à la génération d'un plan d'exécution distribué. Le serveur TiDB est intrinsèquement sans état ; plusieurs instances peuvent être déployées et gérées par un équilibreur de charge (tel que LVS, HAProxy ou F5) pour fournir une adresse d'accès unifiée. Les connexions client sont réparties équitablement entre ces instances. Le serveur TiDB lui-même ne stocke pas de données, il interprète simplement le SQL et transfère les demandes de lecture réelles aux nœuds de stockage sous-jacents, TiKV (ou TiFlash).
  • Serveur PD (Placement Driver) : Agissant comme le module de gestion des métadonnées pour l'ensemble du cluster TiDB, le PD est responsable de la conservation de l'état de distribution des données et de la topologie globale du cluster en temps réel. Il offre également une interface de gestion (TiDB Dashboard) et attribue les identifiants de transaction pour les transactions distribuées. Au-delà du stockage des métadonnées, le PD émet des commandes de planification des données aux nœuds TiKV en fonction des rapports de distribution de données en temps réel, fonctionnant ainsi comme le "cerveau" du cluster. Pour une haute disponibilité, le PD est lui-même constitué d'au moins trois nœuds, et il est recommandé de déployer un nombre impair de nœuds PD.
  • Nœuds de Stockage :
    • Serveur TiKV : Ce composant est en charge du stockage des données. D'un point de vue externe, TiKV fonctionne comme un moteur de stockage clé-valeur distribué et transactionnel. L'unité fondamentale de stockage est le Region, chaque Region gérant une plage de clés (un intervalle semi-ouvert de StartKey à EndKey). Chaque nœud TiKV peut gérer plusieurs Regions. L'API de TiKV offre un support natif pour les transactions distribuées au niveau des paires clé-valeur, fournissant par défaut le niveau d'isolation SI (Snapshot Isolation), ce qui est crucial pour le support des transactions distribuées de TiDB au niveau SQL. Après l'analyse SQL, la couche SQL de TiDB traduit le plan d'exécution en appels effectifs à l'API TiKV. Les données sont donc intégralement stockées dans TiKV. De plus, TiKV maintient automatiquement des réplicas multiples de données (trois par défaut), assurant une haute disponibilité et un basculement automatique.
    • TiFlash : Il s'agit d'un type spécial de nœud de stockage. Contrairement aux nœuds TiKV ordinaires, TiFlash stocke les données sous forme colonnaire. Sa fonction principale est d'accélérer les scénarios d'analyse de données (OLAP).

Mécanismes de Stockage de la Base de Données

Paires Clé-Valeur

Pour un système de persistance de données, le choix du modèle de stockage est primordial. TiKV a opté pour un modèle Clé-Valeur, complété par la capacité de parcourir les clés de manière ordonnée. Deux points essentiels caractérisent le stockage de données de TiKV :

  1. Il s'agit d'une immense structure de type Map (comparable à std::map en C++), qui stocke des paires Clé-Valeur.
  2. Les paires Clé-Valeur de cette Map sont ordonnées selon l'ordre binaire des clés. Cela signifie qu'il est possible de se positionner sur une clé spécifique et de récupérer séquentiellement les paires Clé-Valeur dont les clés sont supérieures en ordre croissant.

Il est important de noter que le modèle de stockage KV de TiKV n'a pas de lien direct avec les tables SQL à ce niveau. Cette discussion se concentre sur l'implémentation d'un stockage Clé-Valeur distribué, performant et fiable, sans aborder les concepts SQL.

Stockage Local (RocksDB)

Toute base de données persistante doit finalement stocker ses données sur disque. TiKV ne fait pas exception. Cependant, au lieu d'écrire directement sur le disque, TiKV confie la persistance des données à RocksDB. RocksDB, un moteur de stockage KV mono-instance open source par Facebook, est une solution extrêmement performante qui répond aux exigences de TiKV. On peut simplifier en considérant RocksDB comme une Map Clé-Valeur persistante et locale.

Le Protocole Raft

Une des plus grandes difficultés rencontrées par TiKV est de garantir l'intégrité et la disponibilité des données en cas de défaillance d'un nœud unique. La solution réside dans la réplication des données sur plusieurs machines. Ainsi, si une machine tombe en panne, les réplicas sur d'autres machines peuvent prendre le relais. Cela nécessite une solution de réplication fiable, efficace et capable de gérer les défaillances de réplicas. TiKV a choisi l'algorithme Raft.

Raft est un protocole de consensus qui garantit la cohérence. Il offre plusieurs fonctionnalités clés :

  • Élection du Leader (réplica primaire)
  • Modification des membres du groupe (ajout/suppression de réplicas, transfert du Leader)
  • Réplication des journaux (logs)

TiKV utilise Raft pour la réplication des données. Chaque modification de données est enregistrée comme une entrée de journal Raft, qui est ensuite répliquée de manière sécurisée et fiable vers chaque nœud du groupe de réplication. Conformément au protocole Raft, l'écriture est considérée comme réussie dès que la synchronisation est effectuée vers une majorité des nœuds.

En résumé, RocksDB permet à TiKV de stocker rapidement les données sur disque, tandis que Raft assure la réplication de ces données sur plusieurs machines pour prévenir les défaillances. Les écritures passent par l'interface Raft, et non directement par RocksDB. Grâce à Raft, TiKV devient un stockage Clé-Valeur distribué et tolérant aux pannes, capable de compléter automatiquement les réplicas manquants en cas de défaillance d'un petit nombre de machines, sans impact perceptible sur les applications.

Regions

Pour faciliter la compréhension, imaginons dans cette section que les données n'ont qu'un seul réplica. TiKV peut être vu comme une immense Map KV ordonnée. Pour permettre une mise à l'échelle horizontale du stockage, les données sont réparties sur plusieurs machines. Il existe deux approches typiques pour distribuer les données dans un système KV :

  • Hachage : Les clés sont hachées, et la valeur de hachage détermine le nœud de stockage correspondant.
  • Plage : Les clés sont divisées en plages contiguës, et chaque plage est stockée sur un nœud spécifique.

TiKV a choisi la deuxième méthode, divisant l'espace de clés en de nombreuses sections, appelées Regions, chacune représentant une série de clés contiguës. Chaque Region est conçue pour ne pas dépasser une certaine taille (actuellement 96 Mo par défaut dans TiKV) et est définie par un intervalle semi-ouvert [StartKey, EndKey).

Ces Regions sont indépendantes des tables SQL. Une fois les données partitionnées en Regions, TiKV effectue deux opérations cruciales :

  • Distribution des données : Les Regions sont réparties sur tous les nœuds du cluster, de manière à équilibrer le nombre de Regions par nœud. Cela permet une scalabilité horizontale de la capacité de stockage (les nouvelles nœuds attirent automatiquement des Regions des autres nœuds) et un équilibrage de charge (évitant que certains nœuds soient surchargés tandis que d'autres sont inactifs). Un composant dédié (le PD) maintient un enregistrement de la distribution des Regions, permettant de localiser n'importe quelle clé.
  • Réplication et gestion par Raft : La réplication et la gestion des membres sont effectuées par Region. Chaque Region possède plusieurs réplicas, appelés Replicas. Ces Replicas sont répartis sur différents nœuds et forment un Groupe Raft. Un des Replicas est désigné comme Leader de ce groupe, les autres étant des Followers. Par défaut, toutes les opérations de lecture et d'écriture sont dirigées vers le Leader ; les lectures sont complétées par le Leader, tandis que les écritures sont répliquées par le Leader vers les Followers.

En distribuant et en répliquant les données par Region, TiKV devient un système Clé-Valeur distribué avec une tolérance aux pannes intégrée, résolvant les problèmes de capacité de stockage et de perte de données due à une défaillance de disque.

MVCC (Multi-Version Concurrency Control)

TiKV implémente le contrôle de concurrence multi-version (MVCC), une fonctionnalité courante dans de nombreuses bases de données. Imaginez deux clients tentant de modifier simultanément la valeur d'une même clé. Sans MVCC, des verrous seraient nécessaires, entraînant des problèmes de performance et de blocage dans un environnement distribué.

L'implémentation MVCC de TiKV est réalisée en ajoutant un numéro de version à la fin de la clé. Sans MVCC, le stockage TiKV pourrait ressembler à ceci :

CléA -> Valeur
CléB -> Valeur
...

Avec MVCC, l'agencement des clés dans TiKV est le suivant :

CléA_Version3 -> Valeur
CléA_Version2 -> Valeur
CléA_Version1 -> Valeur
...
CléB_Version4 -> Valeur
CléB_Version3 -> Valeur
CléB_Version2 -> Valeur
CléB_Version1 -> Valeur
...

Les versions plus récentes d'une même clé sont placées avant les versions plus anciennes. Ce mécanisme garantit que différentes transactions peuvent lire et écrire différentes versions d'une donnée sans se bloquer mutuellement, améliorant ainsi la concurrence.

Transactions Distribuées ACID

Les transactions de TiKV sont basées sur le modèle Percolator, utilisé par Google dans BigTable. TiKV a implémenté ce modèle en y apportant de nombreuses optimisations pour garantir les propriétés ACID (Atomicité, Cohérence, Isolation, Durabilité) dans un environnement distribué.

Couche de Calcul de la Base de Données

Sur la base des capacités de stockage distribué de TiKV, TiDB construit un moteur de calcul capable de gérer efficacement les transactions (OLTP) et d'analyser les données (OLAP). Cette section décrit d'abord comment TiDB mappe les données des schémas relationnels (tables, index) en paires (Clé, Valeur) dans TiKV, puis aborde la gestion des métadonnées de TiDB, et enfin présente l'architecture principale de la couche SQL de TiDB.

Pour la couche de calcul, nous nous concentrerons sur la structure de stockage en lignes basée sur TiKV. Pour les charges de travail analytiques, TiDB propose également TiFlash, une extension de TiKV, qui utilise une structure de stockage en colonnes.

Correspondance entre les Données de Table et les Paires Clé-Valeur

Cette sous-section détaille comment les données SQL sont transformées en paires (Clé, Valeur) pour le stockage dans TiKV. Cela concerne principalement :

  • Les données de chaque ligne d'une table (données de ligne).
  • Les données de tous les index d'une table (données d'index).

Mappage des Données de Ligne

Dans une base de données relationnelle, une table peut contenir de nombreuses colonnes. Pour mapper les données de chaque ligne en une paire (Clé, Valeur), la construction de la clé est cruciale. Pour les scénarios OLTP, qui impliquent un grand nombre d'opérations d'insertion, suppression, mise à jour et interrogation sur des lignes individuelles ou multiples, la base de données doit être capable de lire rapidement une ligne. Par conséquent, la clé doit idéalement contenir un identifiant unique (explicite ou implicite) pour une localisation rapide. De plus, de nombreuses requêtes OLAP nécessitent un balayage complet de la table. Si les clés de toutes les lignes d'une table peuvent être encodées dans une plage contiguë, les balayages complets peuvent être effectués efficacement via des requêtes de plage.

En tenant compte de ces aspects, TiDB conçoit le mappage des données de table aux paires Clé-Valeur comme suit :

  • Pour regrouper les données d'une même table et faciliter leur recherche, TiDB attribue un identifiant unique à chaque table, appelé ID_Table (un entier unique au sein du cluster).
  • TiDB attribue également un identifiant unique à chaque ligne d'une table, appelé ID_Ligne (un entier unique au sein de la table). TiDB optimise ce processus : si une table possède une clé primaire de type entier, la valeur de cette clé primaire est utilisée comme ID_Ligne.

Chaque ligne de données est encodée en une paire (Clé, Valeur) selon la règle suivante :

Clé:   prefixe_table{ID_Table}_separateur_enregistrement{ID_Ligne}
Valeur: [colonne1, colonne2, colonne3, colonne4]

prefixe_table et separateur_enregistrement sont des constantes de chaîne spécifiques utilisées pour distinguer ce type de données des autres dans l'espace des clés. Leurs valeurs exactes seront données plus loin.

Mappage des Données d'Index

TiDB prend en charge les clés primaires et les index secondaires (y compris les index uniques et non uniques). Similairement au mappage des données de table, TiDB attribue un identifiant unique à chaque index d'une table, appelé ID_Index.

Pour les clés primaires et les index uniques, il est nécessaire de localiser rapidement l'ID_Ligne correspondant à partir d'une valeur de clé. L'encodage en paire (Clé, Valeur) suit cette règle :

Clé:   prefixe_table{ID_Table}_separateur_index{ID_Index}_valeursColonnesIndexees
Valeur: ID_Ligne

Pour les index secondaires non uniques, où une valeur de clé peut correspondre à plusieurs lignes, il est nécessaire de rechercher l'ID_Ligne à partir d'une plage de valeurs. L'encodage suit cette règle :

Clé:   prefixe_table{ID_Table}_separateur_index{ID_Index}_valeursColonnesIndexees_{ID_Ligne}
Valeur: null

Synthèse des Relations de Mappage

Les règles d'encodage ci-dessus utilisent les constantes de chaîne suivantes pour prefixe_table, separateur_enregistrement et separateur_index, qui servent à différencier les types de données dans l'espace des clés :

prefixe_table       = []octet{'t'}
separateur_enregistrement = []octet{'r'}
separateur_index    = []octet{'i'}

Il est important de noter que, dans toutes ces schémas d'encodage, toutes les lignes d'une table partagent le même préfixe de clé, et toutes les données d'un index partagent également un préfixe commun. Les données ayant le même préfixe sont ainsi regroupées dans l'espace des clés de TiKV. En concevant soigneusement l'encodage de la partie suffixe pour maintenir l'ordre de comparaison avant et après l'encodage, les données de table ou d'index peuvent être stockées de manière ordonnée dans TiKV. Grâce à cet encodage, toutes les lignes d'une table sont ordonnées par ID_Ligne dans l'espace des clés de TiKV, et les données d'un index sont ordonnées selon les valeurs spécifiques des colonnes indexées.

Exemple de Mappage Clé-Valeur

Prenons un exemple simple pour illustrer le mappage Clé-Valeur de TiDB. Supposons la table suivante :

CREATE TABLE Produit (
    ProduitID INT,
    Nom VARCHAR(50),
    Categorie VARCHAR(30),
    Prix DECIMAL(10, 2),
    PRIMARY KEY (ProduitID),
    KEY idx_categorie (Categorie)
);

Et trois lignes de données :

1, "Laptop", "Électronique", 1200.00
2, "Clavier", "Accessoire", 75.50
3, "Souris", "Accessoire", 25.00

Puisque la table a une clé primaire de type entier, ProduitID est utilisé comme ID_Ligne. Si l'ID_Table est 20, les données de table stockées dans TiKV seraient :

t20_r1 --> ["Laptop", "Électronique", 1200.00]
t20_r2 --> ["Clavier", "Accessoire", 75.50]
t20_r3 --> ["Souris", "Accessoire", 25.00]

La table a également un index secondaire non unique idx_categorie. Si l'ID_Index de cet index est 5, les données d'index stockées dans TiKV seraient :

t20_i5_Accessoire_2     --> null  (pour "Clavier")
t20_i5_Accessoire_3     --> null  (pour "Souris")
t20_i5_Électronique_1   --> null  (pour "Laptop")

Cet exemple illustre les règles de mappage du modèle relationnel de TiDB vers le modèle Clé-Valeur, ainsi que les considérations qui ont guidé ces choix.

Gestion des Métadonnées

Chaque Base_de_données et Table dans TiDB possède des métadonnées, y compris sa définition et ses propriétés. Ces informations doivent également être persistantes et sont stockées dans TiKV. Chaque Base_de_données/Table reçoit un identifiant unique, qui est encodé dans la clé des paires Clé-Valeur de métadonnées, précédé du préfixe m_. La valeur stockée est la métadonnée sérialisée.

De plus, TiDB utilise une paire (Clé, Valeur) dédiée pour stocker le numéro de version le plus récent de toutes les informations de schéma de table. Cette paire est globale et sa version est incrémentée à chaque modification de l'état d'une opération DDL. Actuellement, cette paire est stockée de manière persistante dans le serveur PD, avec la clé "/tidb/ddl/global_schema_version" et une valeur de type int64. TiDB s'inspire de l'algorithme de modification de schéma en ligne de Google F1, avec un thread d'arrière-plan vérifiant constamment les changements de version du schéma dans le serveur PD, garantissant que les modifications de version sont détectées dans un laps de temps défini.

Introduction à la Couche SQL

La couche SQL de TiDB, incarnée par le Serveur TiDB, est chargée de traduire les requêtes SQL en opérations Clé-Valeur. Ces opérations sont ensuite transmises à la couche de stockage distribuée partagée TiKV. Le Serveur TiDB assemble ensuite les résultats renvoyés par TiKV et les retourne au client. Les nœuds de cette couche sont sans état et ne stockent pas de données, étant tous égaux.

Optimisation des Opérations SQL

Une approche simple de la couche de calcul consisterait à traduire chaque opération SQL en une série d'accès individuels à TiKV. Cependant, cela peut entraîner de nombreuses appels RPC et une inefficacité, en particulier pour les opérations complexes.

Pour résoudre ce problème, les opérations de calcul doivent être exécutées aussi près que possible des nœuds de stockage afin de minimiser les appels RPC. Par exemple, une condition de prédicat SQL comme nom = "TiDB" devrait être "poussée" vers les nœuds de stockage pour y être évaluée. Cela permet de ne renvoyer que les lignes pertinentes, évitant ainsi des transferts réseau inutiles. De même, les fonctions d'agrégation comme COUNT(*) peuvent être "pré-agrégées" au niveau des nœuds de stockage ; chaque nœud renverrait un résultat partiel de COUNT(*), que la couche SQL additionnerait ensuite pour obtenir le total.

Ce modèle réduit considérablement le volume de données transférées sur le réseau et la charge de traitement de la couche SQL.

Architecture de la Couche SQL

Les requêtes SQL des utilisateurs sont envoyées directement ou via un équilibreur de charge aux Serveurs TiDB. Le Serveur TiDB analyse le paquet du protocole MySQL, extrait le contenu de la requête, procède à l'analyse syntaxqiue et sémantique du SQL, élabore et optimise un plan de requête, puis exécute ce plan pour récupérer et traiter les données. Étant donné que toutes les données sont stockées dans le cluster TiKV, le Serveur TiDB interagit avec TiKV pendant ce processus pour obtenir les données. Enfin, le Serveur TiDB renvoie les résultats de la requête à l'utilisateur.

Étiquettes: TiDB Architecture Distribuée TiKV raft SQL

Publié le 19 juillet à 23h54