Initiation à l’algorithmique répartie

Wiley
Page 1 sur 95Lecteur de document UniversityLib

Initiation à l’algorithmique répartie

Algorithmique, Systèmes répartis · course

Voir tous les documents en programmation

Initiation à l’algorithmique répartie

Denis Conan

Revision : 151

CSC4509

Télécom SudParis

Avril 2020

Initiation à l’algorithmique répartie

Table des matières

Initiation à l’algorithmique répartie

Denis Conan, , Télécom SudParis, CSC4509 Avril 2020

Licence

Utilisation du cours

Plan du document

1 Éléments introductifs

1.1 Modèle de système réparti . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.1.1 Modèle de transitions 1.1.2 Synchronisme . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.1.3 Types de défaillances . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.2 Conventions de codage des algorithmes répartis . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.3 Relation « arrivé avant » aussi appelée précédence causale, Lamport 1978 . . . . . . . . . . . . 1.3.1 Algorithme de calcul des horloges scalaires de Lamport 1978 . . . . . . . . . . . . . . . 1.3.2 Algorithme de calcul des horloges vectorielles de Fidge 1991 . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.3.3 Exercices 1.4 Vague et traversée de graphe . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.4.1 Algorithme de vague centralisé Écho de Segall, 1983 . . . . . . . . . . . . . . . . . . . . 1.4.2 Algorithme de vague décentralisé de Finn, 1979 * . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1.4.3 Exercices

2 Élection

2.1 Propriétés et vocabulaire . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 2.2 Élection dans un anneau, algorithme de Le Lann, 1977 . . . . . . . . . . . . . . . . . . . . . . . 2.3 Élection avec l’algorithme de vague Écho de Segall, 1983 . . . . . . . . . . . . . . . . . . . . . . 2.3.1 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .

3 Diffusion

3.1 Spécification des diffusions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.2 Diffusion fiable . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.2.1 Algorithme de diffusion fiable . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.3 Diffusion FIFO . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.3.1 Algorithme de diffusion FIFO . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.4 Diffusion causale . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.4.1 Algorithme de diffusion causale construit à partir d’un algorithme de diffusion FIFO . . 3.4.2 Algorithme de diffusion causale à base d’horloge vectorielle de Birman et Joseph, 1987 . 3.4.3 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.5 Diffusion atomique (ou totale) 3.6 Relations entre les diffusions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.7 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.8 Diffusion atomique et consensus * . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.8.1 Résultat d’impossibilité du consensus * . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.8.2 Algorithmes de diffusion temporisée * . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.9 Propriété d’uniformité * . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3.10 Inconsistance et contamination * . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .

4 Exclusion mutuelle

4.1 Propriétés . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 4.2 Algorithmes à base de permissions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 4.2.1 Structure informationnelle générique de Sanders, 1987 . . . . . . . . . . . . . . . . . . . 4.2.2 Algorithme générique de Sanders, 1987 . . . . . . . . . . . . . . . . . . . . . . . . . . . . 4.2.3 Quelques algorithmes (dérivés de l’algorithme générique) . . . . . . . . . . . . . . . . . .

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

1

4

5

6

7 8 9 11 12 14 15 16 17 18 19 20 22 23

24 25 26 28 30

31 32 33 34 35 36 38 40 42 44 45 46 47 48 49 50 52 53

54 55 56 57 58 59

2

Initiation à l’algorithmique répartie

4.3 Algorithme à base de jeton de Ricart et Agrawala 1983, et de Suzuki et Kasami, 1985 . . . . . 4.4 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .

62 64

5 Interblocage

65 5.1 Principaux modèles d’interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 66 5.1.1 Modèle d’interblocage ET . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 67 5.1.2 Modèle d’interblocage OU–ET . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 68 5.2 Condition de déblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 69 5.3 Définition de l’interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 70 5.4 Trois stratégies contre l’interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 71 5.5 Prévention dans le modèle ET avec l’algorithme de Rosenkrantz, Stearns et Lewis, 1978 . . . . 73 5.5.1 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 74 5.6 Détection d’interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 75 76 5.6.1 Coupure cohérente . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 5.6.2 Algorithme centralisé de construction de coupure cohérente de Chandy et Lamport, 1985 77 78 5.6.3 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .

6 Détection de terminaison

6.1 Modèle OU de l’interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.2 Configurations terminale et finale . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.3 États actif et passif, et algorithme de contrôle . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.4 Algorithmes de détection de terminaison . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.5 Détection par calcul du graphe d’exécution . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.5.1 Algorithme de Dijkstra et Scholten, 1980 . . . . . . . . . . . . . . . . . . . . . . . . . . 6.5.2 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.6 Détection par vagues dans un anneau . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.6.1 Algorithme de Safra, 1987 . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6.6.2 Exercice . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .

Bibliographie

Index

Fin

79 80 81 82 83 84 85 87 88 89 91

92

93

95

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

3

Initiation à l’algorithmique répartie

’

Licence

$

Ce document est une documentation libre, placée sous la Licence de Documentation Libre GNU (GNU Free Documentation License).

Copyright (c) 2003-2020 Denis Conan Permission est accordée de copier, distribuer et/ou modifier ce document selon les termes de la Licence de Documentation Libre GNU (GNU Free Documentation License), version 1.2 ou toute version ultérieure publiée par la Free Software Foundation; avec les Sections Invariables qui sont ‘Licence’ ; avec les Textes de Première de Couverture qui sont ‘Initiation à l’algorithmique répartie’ et avec les Textes de Quatrième de Couverture qui sont ‘Fin’. Une copie de la présente Licence peut être trouvée à l’adresse suivante : http://www.gnu.org/copyleft/fdl.html.

2

Remarque : La licence comporte notamment les sections suivantes : 2. COPIES VERBATIM, 3. COPIES

EN QUANTITÉ, 4. MODIFICATIONS, 5. MÉLANGE DE DOCUMENTS, 6. RECUEILS DE

DOCUMENTS, 7. AGRÉGATION AVEC DES TRAVAUX INDÉPENDANTS et 8. TRADUCTION.

&

%

Ce document est préparé avec des logiciels libres :

• LATEX :

les

textes

sources

sont écrits en LATEX (http://www.latex-project.org/,

le site du Groupe francophone des Utilisateurs de TEX/LATEX est http://www.gutenberg.eu.org). la classe seminar ont Une nouvelle été fusionforge slideint, tout https://fusionforge.int-evry.fr/www/slideint/);

spécialement dévéloppées: newslide et slideint (projet

et une nouvelle

style basées

feuille de

classe

sur

• emacs: tous les textes sont édités avec l’éditeur GNU emacs (http://www.gnu.org/software/emacs); • dvips: les versions PostScript (PostScript est une marque déposée de la société Adobe Systems In- corporated) des transparents et des polycopiés à destination des étudiants ou des enseignants sont obtenues à partir des fichiers DVI (« DeVice Independent ») générés à partir de LaTeX par l’utilitaire dvips (http://www.ctan.org/tex-archive/dviware/dvips);

• ps2pdf et dvipdfmx: les versions PDF (PDF est une marque déposée de la société Adobe Sys- tems Incorporated) sont obtenues à partir des fichiers Postscript par l’utilitaire ps2pdf (ps2pdf étant un shell-script lançant Ghostscript, voyez le site de GNU Ghostscript http://www.gnu.org/- software/ghostscript/) ou à partir des fichiers DVI par l’utilitaire dvipfmx;

• makeindex:

les

index et glossaire

Publicité

sont générés à l’aide de

l’utilitaire Unix makeindex

(http://www.ctan.org/tex-archive/indexing/makeindex);

• TeX4ht: les pages HTML sont générées à partir de LaTeX par TeX4ht (http://www.cis.ohio-

-state.edu/~gurari/TeX4ht/mn.html);

• Xfig: les figures sont dessinées dans l’utilitaire X11 de Fig xfig (http://www.xfig.org); • fig2dev: les figures sont exportées dans les formats EPS (« Encapsulated PostScript ») et PNG (« Portable Network Graphics ») grâce à l’utilitaire fig2dev (http://www.xfig.org/userman/- installation.html);

• convert: certaines figures sont converties d’un format vers un autre par l’utilitaire convert

(http://www.imagemagick.org/www/utilities.html) de ImageMagick Studio;

• HTML TIDY:

les sources HTML générés par TeX4ht sont « beautifiés » à l’aide de HTML TIDY

(http://tidy.sourceforge.net) ; vous pouvez donc les lire dans le source.

Nous espérons que vous regardez cette page avec un navigateur libre: Firefox par exemple. Comme l’indique le choix de la licence GNU/FDL, tous les éléments permettant d’obtenir ces supports sont libres. Ce cours a bénéficié des relectures attentives et constructives de François Meunier, Léon Lim.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

4

Initiation à l’algorithmique répartie

’

Utilisation du cours

$

■ Apprentissage en formation en ligne

♦ Démarche conseillée :

3

▶ Pour chaque page du support de cours, étudiez la page du support de cours

(diapositive + commentaires)

■ Étude des exercices en présentiel ■ Puis auto-évaluation des connaissances avec les QCM (une série par section) ■ Les QCM ainsi que les corrigés des exercices sont fournis à part dans moodle

&

%

NB : Certaines pages du cours sont marquées par un astérisque (« * ») à la fin de leur titre. Ceci correspond à un contenu d’approfondissement. Ne l’étudiez pas en détail avant de maîtriser les autres points de la section. Ces diapositives doivent être étudiées avant de faire certaines questions (optionnelles) des exercices.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

5

Initiation à l’algorithmique répartie

’

Plan du document

$

4

1 Éléments introductifs . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 5 2 Élection . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 19 3 Diffusion . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 24 4 Exclusion mutuelle . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 42 5 Interblocage . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 50 6 Détection de terminaison . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 63

&

%

Une définition communément admise d’un système réparti est : « plusieurs ordinateurs inter-connectés par un ensemble de réseaux de communication effectuant un travail ensemble et communicant par échange de messages ». Cette définition suggère deux propriétés principales des systèmes répartis : la non-unicité de lieu (un utilisateur travaille en local ou à distance) et la non-unicité de temps (chaque ordinateur possède sa propre notion du temps à travers son horloge physique). Par ailleurs, cette définition diffère de celle d’un système parallèle (ou dit fortement couplé) dans lequel les communications peuvent s’effectuer via une mémoire partagée et dans lequel il y a unicité de lieu et de temps.

Les systèmes répartis sont difficiles à concevoir et à comprendre parce qu’ils ne sont pas intuitifs. Peut- être est-ce aussi parce que notre vie est, par bien des aspects, fondamentalement séquentielle ? Nous devons donc développer une intuition pour la répartition. Dans cette discipline, il existe une tension inévitable entre les partisans de la modélisation et de l’analyse, et ceux de l’observation expérimentale. Cette tension illustre la dichotomie classique entre la théorie et la pratique. Dans ce cours d’algorithmique répartie, nous nous placerons plus du côté pratique.

Nous commençons par la présentation du modèle de système réparti dans la section introductive. Ensuite, les problèmes étudiés sont des problèmes fondamentaux basiques de l’algorithmique répartie à partir desquels sont construites des architectures de services répartis complexes. Le principe de l’élection est de partir d’une configuration dans laquelle tous les processus sont dans le même état, pour arriver dans une configuration dans laquelle un seul processus est dans l’état « gagnant » et tous les autres dans l’état « perdant ». La diffusion est une primitive de communication permettant à un processus d’envoyer le même message à tous les autres processus en respectant des propriétés d’ordre (FIFO, causal, total) dans la transmission des messages. L’exclusion mutuelle consiste à faire circuler un jeton entre des processus répartis sur le réseau pour n’autoriser qu’un seul d’entre eux à entrer en section critique. C’est l’expression répartie des sémaphores. L’interblocage se produit lorsqu’un ensemble de processus est tel que chacun d’eux tient au moins une ressource, et pour poursuivre sa progression, est en attente d’une ressource tenue par l’un des autres. Des algorithmes sont proposés pour prévenir ou détecter les interblocages. La détection de terminaison autorise qu’un algorithme réparti se termine de façon implicite, c’est-à-dire sans que les processus atteignent leur état final (fin du processus par appel de la fonction exit), autrement dit, « parce qu’il n’y a plus de travail à faire » (tous les processus sont en attente d’un message et les canaux de communication sont vides).

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

6

Initiation à l’algorithmique répartie

’

$

1 Éléments introductifs

5

1.1 Modèle de système réparti . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6 1.2 Conventions de codage des algorithmes répartis . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 10 1.3 Relation « arrivé avant » aussi appelée précédence causale, Lamport 1978 . . . . . . . 11 1.4 Vague et traversée de graphe . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 15

&

%

Cette section introductive est reprise des références suivantes :

• G. Tel, Chapter 2 : The Model, dans Introduction to Distributed Algorithms, Cambridge University

Press, pp. 43–72, 1994.

• C. Fidge. Logical Time in Distributed Computing Systems, dans IEEE Computer, pages 28–33, August

1991.

• G. Tel, Chapter 6 : Wave and traversal algorithms, dans Introduction to Distributed Algorithms, Cam-

bridge University Press, pp. 177–221, 1994.

• H. Attiya et J. Welch, Chapter 2 : Basic Algorithms in Message-Passing Systems, dans Distributed

Computing : Fundamentals, simulation, and advanced topics, Wiley, pp. 9–30, 2004.

• J.H. Saltzer, M.F. Kaashoek, Principles of Computer System Design : An Introduction, Morgan Kauf-

mann, 2009.

Le premier élément d’introduction est le modèle de système réparti. Ce modèle à base de messages est utilisé dans tout le reste du cours. Il est assez général pour être utile aussi bien lors de la conception que lors de la vérification (même si nous ne nous focalisons par sur les preuves des algorithmes étudiés). Puis, nous introduisons les conventions de codage utilisées dans ce cours : soit l’orientation contrôle soit l’orientation évènement. Ensuite, le dernier élément général introduit pour la suite du cours est la notion de dépendance causale qui permet de construire un ordre partiel des évènements d’une exécution répartie. Enfin, parmi les problèmes fondamentaux que nous étudions, beaucoup peuvent s’exprimer à l’aide de sous-tâches génériques comme les vagues. C’est par exemple le cas de la diffusion ou de la détection d’interblocage étudiée un peu plus loin dans ce cours. Nous présentons donc à la fin de cette section le principe des algorithmes de vagues et y ferons référence dans les autres sections.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

7

Initiation à l’algorithmique répartie

’

1 Éléments introductifs $

1.1 Modèle de système réparti

6

1.1.1 Modèle de transitions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 7 1.1.2 Synchronisme . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 8 1.1.3 Types de défaillances . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 9

&

%

Les paramètres du système réparti les plus discriminants sont les suivants : synchronisme (« performance » des nœuds et du réseau), type de fautes des nœuds et des communications (catégorie de fautes matérielles et logicielles [persistantes ̸= transitoires]), topologie (forme du graphe du système de communication) et déterminismes des processus (comportement prévisible des algorithmes répartis). Ces hypothèses sont im- portantes car elles conditionnent les résultats d’impossibilité (comme l’atteinte d’un consensus) ainsi que les approches algorithmiques. Par exemple, il n’existe pas de solution déterministe à tous les problèmes dans tous les cas. Mais, avant de détailler ces éléments discriminants, nous présentons les éléments constitutifs du modèle, c’est-à-dire dans notre cas ceux du modèle de transitions.

Ce cours se limite à l’étude des algorithmes répartis déterministes. Ainsi, nous ne présentons pas d’algo- rithme avec comportement aléatoire, par exemple ceux du type « Las Vegas » (avec une solution correcte mais une distribution probabiliste sur la durée d’exécution) ou ceux du type « Monte Carlo » (avec une probabilité sur la correction de la solution mais un bornage de cette probabilité).

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

8

1 Éléments introductifs

Publicité

’

1.1 Modèle de système réparti $

1.1.1 Modèle de transitions

7

■ Algorithme : assigne des valeurs à des variables ■ Processus : exécution d’un algorithme sur un nœud du système réparti ■ Canal de communication : lien logique réseau entre deux processus ■ Exécution = état initial puis succession d’actions sur les états ■ État global ou configuration = juxtaposition d’états locaux, un par processus ■ Histoire = séquence des actions d’un processus

♦ Histoire répartie = ensemble des histoires locales des processus

■ Répartition =⇒ actions internes, d’émission et de réception ■ Toute propriété est ou se décompose en propriétés de :

♦ Correction (en anglais, safety ) : une assertion est vraie dans chaque

configuration de l’algorithme

♦ Vivacité ou progression (en anglais, liveness) : une assertion est vraie dans

certaines configurations de chaque exécution

&

%

Tous les modèles de spécification pour les systèmes répartis sont basés sur la notion d’action atomique et de machine à états. Parmi les modèles les plus couramment choisis, citons le modèle des transitions, le modèle des actions temporelles logiques, et le modèle des automates (pouvant être temporisés). Dans ce cours, nous utilisons le modèle le plus simple, celui dit des transitions, que nous décrivons maintenant.

Un algorithme assigne des valeurs à des variables. Un état est l’affectation de valeurs à des variables. Une action (encore appelée un évènement) a représente la relation entre un ancien état s et un nouvel état t notée sat. Les actions sont atomiques : l’action a provoque le changement « instantané » de l’état de s à t. Autrement dit, on n’observe pas l’état du système pendant l’exécution de a.

Un algorithme séquentiel A est un objet syntaxique construit selon la grammaire d’un langage de pro- grammation. Une exécution de l’algorithme séquentiel A dans un processus P débutant à l’instant t à partir de l’état initial s0, est la succession d’un nombre infini d’actions sur des états notée a0s0a1s1a2s2a3, etc. L’état s0 est appelé l’état initial. L’action initiale a0 est fictive et correspond à la création du processus P . L’exécution d’un algorithme séquentiel produit donc une séquence d’actions. La séquence d’actions est appelée l’histoire de l’algorithme séquentiel. De façon duale, lorsqu’il s’agit de formuler ou de démontrer une propriété, il est souvent très intéressant de modéliser une exécution comme étant une séquence d’états, une action étant la transition d’un état à un autre 1,2.

Un algorithme réparti A est composé d’algorithmes séquentiels exécutés dans des processus P1...Pn qui communiquent par échange de messages. L’action d’émission émettre(Pj, m) d’un message m de Pi vers Pj ajoute m au canal cij. Pratiquement, m est transmis de Pi vers Pj par le réseau de communication et est gardé dans la mémoire du nœud où s’exécute Pj. L’action de réception recevoir(m) d’un message m par Pj sur l’ensemble des canaux cij assigne à m le premier message arrivé par l’un des canaux. Si aucun message n’est arrivé alors Pj attend jusqu’à l’arrivée d’un message par l’un des canaux. Toute action autre qu’une action d’émission ou de réception d’un message est appelée une action interne.

L’état d’un canal cij à l’instant physique t est constitué de l’ensemble des messages émis et non encore reçus. L’état global (encore appelé configuration) s d’un système réparti à l’instant physique t est composé des états locaux de tous les processus du système réparti à l’instant t. L’exécution d’un algorithme réparti débutant à l’instant physique t à partir de la configuration s est constituée de la juxtaposition des exécutions des processus. Par déduction, l’histoire répartie d’un système réparti est constituée de la juxtaposition des histoires des processus.

1. R.W. Floyd. Assigning meanings to programs. In J.T. Schwartz, editor, Proceedings of Symposia in Applied Mathema- tics, Mathematical Aspects of Computer Science, volume 19, pages 19–32, Providence, Rhode Island, USA, 1967. American Mathematical Society.

2. C.A.R. Hoare. An Axiomatic Basis for Computer Programming. Communications of the ACM, 12(10):576–580, October

1969.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

9

1 Éléments introductifs

1.1 Modèle de système réparti

Le modèle que nous venons de définir est appelé le modèle de transitions sous l’hypothèse de communi- cation asynchrone : les opérations d’émission ne sont pas bloquantes alors que les opérations de réception le sont. Dans le cas des communications synchrones, les opérations d’émission sont elles-aussi bloquantes. Dans ce cas, clairement, dans toutes les configurations dans lesquelles tous les processus viennent d’exécuter une action interne, les canaux sont vides.

L’exécution d’un algorithme réparti est classiquement représentée par un diagramme de séquences (à la UML) aussi appelé diagramme temporel ou chronogramme. La figure qui suit trace un tel diagramme pour un algorithme réparti composé de trois algorithmes séquentiels. Le déroulement du temps est décrit par une ligne continue (ligne qui est à tiret et s’appelle « ligne de vie » en UML) pour chaque algorithme séquentiel. Les actions sont symbolisées par des tirets sur les « lignes de vie ». Les messages sont matérialisés par des flèches connectant une action émettre à une action recevoir.

Enfin, dans un système réparti, il est important de faire la distinction entre les propriétés de sûreté ou correction (en anglais, safety property) et de progression ou vivacité (en anglais, liveness property). La propriété de sûreté d’un algorithme est de la forme « l’assertion est vraie dans chaque configuration de l’algorithme », ou encore de façon informelle « l’assertion est toujours vraie ». Pratiquement, la propriété de sûreté sert à exprimer que quelque chose de non désiré n’arrive pas. La technique de base pour montrer que l’assertion est toujours vraie est de démontrer que c’est un invariant : vrai dans l’état de départ de l’exécution, et si vrai dans la configuration atteignable s alors vrai dans toutes les configurations atteignables directement à partir de s. La propriété de vivacité d’un algorithme quant à elle stipule que « l’assertion est vraie dans certaines configurations de chaque exécution de l’algorithme », ou encore que « l’assertion est vraie à terme ou ultimement ». La technique de base pour montrer que l’assertion est vraie à terme est soit d’utiliser une autre propriété de vivacité (par exemple, le message est reçu à terme par un processus atteignable car tous les canaux de communication du système transmettent in fine tous les messages émis par l’émetteur du canal), soit par induction en utilisant une métrique qui progresse dans le temps jusqu’à atteindre un seuil auquel l’assertion est vraie (pour le même exemple, parmi chaque pas d’exécution considérant l’émission ou la réception d’un message, de temps en temps, il y a un message qui s’approche de récepteur en récepteur du destinataire final, la métrique utilisée étant le nombre [fini] de processus entre l’émetteur initial et le destinataire final).

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

10

222233331230123023aaaaaaaaa42a25a4a26a35a3a276a2a38713Temps1716151121101PPP134aaaaaaaa1 Éléments introductifs

’

1.1 Modèle de système réparti $

1.1.2 Synchronisme

8

■ Synchronea :

∧ Durée de transmission d’un message d’un nœud à un autre bornée et borne

connue

∧ Durée d’éxécution d’une action interne d’un processus bornée et borne connue

■ Asynchrone :

∨ Pas de borne ou borne inconnue sur la transmission d’un message ∨ Pas de borne ou borne inconnue sur la dérive des horloges ∨ Pas de borne ou borne inconnue sur la durée d’un traitement

a. Il est important de ne pas confondre synchronisme des communications et synchronisme du

système. Ici, c’est le synchronisme du système qui est considéré.

&

%

Un système réparti est dit « synchrone » si et seulement si :

• la durée de transmission d’un message d’un nœud à un autre est bornée et la borne est connue ; et, • la durée d’exécution d’une action interne d’un processus est bornée et la borne est connue.

Dans un tel système réparti, un processus émettant un message peut faire l’hypothèse qu’il est reçu, voire

traité, après une durée limite calculable (car les bornes sont connues).

Un système réparti est dit « asynchrone » s’il n’existe pas de borne (connue ou non) sur la transmission d’un message, la dérive des horloges ou la durée d’exécution d’une action interne. Dans la pratique, construire un algorithme pour un système asynchrone signifie ne pas s’occuper des caractéristiques matérielles (qualité des nœuds ou des communications). En d’autres termes, dès que des aspects temporels sont introduits dans les algorithmes (hypothèse sur les durées d’exécution ou de transmission, test de fiabilité de transmission à l’aide de temporisation, etc.), le système considéré n’est plus « complètement asynchrone », mais dit « partiellement asynchrone ». L’acception « partiellement asynchrone » recouvre le modèle « synchrone » et trente-et-un autres modèles, tous entre « complètement asynchrone » et « synchrone »1. Nous ne détaillons pas ces nombreux modèles et n’abordons pas la tolérance aux fautes dans ce manuscrit. C’est l’objectif des études d’articles réalisées par groupe dans le cadre du module.

1. D. Dolev, C. Dwork, and L. Stockmeyer, On the minimal synchronism needed for distributed consensus, Journal of the

ACM, 34(1), January 1987.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

11

1 Éléments introductifs

’

1.1 Modèle de système réparti $

1.1.3 Types de défaillances

9

■ Hormis lorsque précisé, les algorithmes présentés ne tolérent pas les défaillances ■ Dans tous les cas, les défaillances arbitraires sont exclues de l’étude

&

%

Un processus ou un nœud (dans le cas où l’on ne considère qu’un processus par nœud) est dit « défaillant » lors d’une exécution si son comportement diffère de la spécification de l’algorithme qu’il exécute. Sinon, il est dit « correct ». Il en est de même pour les canaux de communication entre processus. Un modèle de défaillances définit les défaillances rencontrées (et prises en compte ou tolérées). La littérature liste communément les types de défaillances suivants :

• arrêt franc initial (en anglais, initial crash) : un processus n’exécute aucune action de son algorithme local, autre que l’action initiale. Ce type de défaillance n’est pas montré sur la figure ; il se trouverait avant le nœud « arrêt franc » dans le graphe ;

• arrêt franc (en anglais, crash) : un processus s’arrête prématurément et ne fait rien ensuite ; avant l’arrêt, son exécution est correcte. Dans le cas d’un canal de communication, celui-ci est définitivement coupé ;

• « omission sur émission » : un processus s’arrête prématurément, omet d’émettre des messages par intermittence ou les deux. Dans le cas d’un canal de communication, l’omission sur émission correspond à une perte de messages. Par exemple, un processus jetant des messages parce que son cache de messages en émission est plein, subit des omissions sur émission ;

• « omission sur réception » : un processus s’arrête prématurément, omet de recevoir des messages par intermittence ou les deux. Dans le cas d’un canal de communication, l’omission sur réception correspond à une perte de messages. Par exemple, un processus jetant des messages parce que son cache de messages en réception est plein, subit des omissions sur réception ;

• « omission générale » : un processus est sujet à omission sur émission ou sur réception, voire les deux ; • « arbitraire, byzantine1 ou maligne » : un processus peut être sujet à n’importe quel comportement, y compris de la malveillance (de la part d’un utilisateur) ; par opposition, les défaillances précédentes sont dites « bénignes ». Pour un canal de communication, cela correspond à la perte, la duplication, la corruption (violation de l’intégrité), voire la génération spontanée d’un message.

Ces types de défaillances peuvent être classés en termes de sévérité. La figure ordonne les types de défaillances, des défaillances les moins sévères (arrêt francs) au plus sévères (arbitraires). Un algorithme tolérant les défaillances arbitraires tolère aussi les arrêts francs.

Les types de défaillances présentés ici existent aussi bien dans les systèmes synchrones qu’asynchrones. En outre, dans les systèmes synchrones, les défaillances peuvent aussi être « temporelles ». Un processus sujet à des défaillances temporelles peut défaillir des manières suivantes :

• omission générale ;

1. Le qualificatif « byzantin » vient de l’article célèbre de Lamport, Shostak et Pease « The Byzantine Generals Problem »

de 1982, ACM Transactions on Programming Languages and Systems, 4(3):382–401, July 1982.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

12

Arrêt francOmission sur émissionOmission sur réceptionOmission généraleArbitraire, byzantine ou maligneplus sévèremoins sévère1 Éléments introductifs

1.1 Modèle de système réparti

Publicité

• « défaillance de l’horloge locale » : l’horloge locale dérive au delà de la borne tolérée ; • « défaillance de performance » : la durée d’exécution d’un traitement dépasse la borne tolérée ou est trop courte. Pour un canal de communication, il s’agit d’une transmission trop rapide ou trop lente.

Dans ce cours, hormis lorsque nous le précisons explicitement dans de rares cas, les algorithmes présentés sont conçus dans le cas de systèmes répartis sans défaillance. En outre, les défaillances arbitraires étant très difficiles à tolérer, le cours ne les aborde pas du tout. La tolérance aux fautes bénignes est étudiée dans les études bibliographiques du module.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

13

Initiation à l’algorithmique répartie

’

1 Éléments introductifs $

1.2 Conventions de codage des algorithmes répartis

■ Notation orientée contrôle

■ Notation orientée évènement

10

1

2

3

4

5

6

7

8

9

10

11

Chaque processus p exécute :

var r[q] pour tout q ∈ V oisins init F ; begin

while #{q : r[q] = F } > 1 do

recevoir⟨jeton⟩ de q r[q] := T

émettre⟨jeton⟩ vers q0 avec r[q0] = F recevoir⟨jeton⟩ de q0 r[q0] := T decider

end

1

2

3

4

5

6

7

8

9

10

11

Chaque processus p exécute :

var r[q] pour tout q ∈ V oisins init F ; var sp, dp init F ; Sp : {#{q : r[q] = F } = 1 and sp = F }

émettre⟨jeton⟩ vers q0 avec r[q0] := F sp := T

Rp : {Un message ⟨jeton⟩ est arrivé}

recevoir⟨jeton⟩ de q r[q] := T

Dp : {#{q : r[q] = F } = 0 et dp = F }

decider ; dp := T

&

%

Dans les deux formes, les opérations de réception recevoir(m) ne spécifient pas le processus émetteur du message, mais l’émetteur est connu après la réception, c’est-à-dire dans la portion d’algorithme traitant la réception. La notation « #E » est utilisée pour signifier la cardinalité de l’ensemble E. Les notations « T » et « F » sont utilisées pour signifier les valeurs booléenes « vrai » et « faux », respectivement. Les deux algorithmes de cette page réalisent le même traitement, mais ce n’est pas ce qui nous intéresse dans cette diapositive.

Les deux orientations (contrôle et évènement) sont possibles pour chaque algorithme, mais dans de nom- breux cas, l’une est plus commode que l’autre. L’orientation contrôle d’un algorithme consiste en un algo- rithme séquentiel par processus avec des actions d’émission et de réception. La structure de l’algorithme est ainsi exprimée explicitement. En revanche, le non-déterminisme est plus facilement exprimé (parce qu’impli- cite) dans l’orientation évènement. La spécification consiste en une déclaration de variables suivie d’une liste d’actions. Chaque action consiste en une expression booléenne (l’action de garde) et un bloc d’instructions correspondantes. L’action est autorisée ou applicable lorsque la garde est évaluée à vrai. Dans ce cas, les instructions de l’action sont exécutées atomiquement, c’est-à-dire sans interruption dans le bloc. Les actions sont exécutées dans n’importe quel ordre de façon non déterministe.

Télécom SudParis — Denis Conan — Avril 2020 — CSC4509

14

Initiation à l’algorithmique répartie

’

1 Éléments introductifs $

1.3 Relation « arrivé avant » aussi appelée précédence causale, Lamport 1978

■ Précédence causale (en anglais, happened before relation) entre évènements :

♦ La relation « → » sur les évènements d’un système est la plus petite relation

satisfaisant les conditions suivantes : ▶ Si a et b sont des évènements sur un même processus, et a arrive avant b,

alors a → b

▶ Si a est l’émission d’un message par un processus et b la réception du même

message par un autre processus, alors a → b

▶ Si a → b et b → c, alors a → c

■ La relation « → » définit un ordre partiel irréflexif ■ a||b = a ̸→ b ∧ b ̸→ a signifie que les deux actions sont concurrentes ■ Ensuite, l’horloge logique prend différentes formes

♦ Scalaire, vecteur ou matrice + avec estampillage des messages ♦ Et lorsque les dépendances causales sont construites à base de vecteurs

11

▶ Dépendances directes : estampilles des messages contenant un scalaire ▶ Dépendances indirectes : estampilles des messages contenant un vecteur

&

%

La non-unicité de temps impliquant le fait qu’il est impossible de réaliser un observateur qui voit tout le système de manière immédiate, empêche l’utilisation d’une horloge physique exacte pour ordonner les actions d’un algorithme réparti. Nous avons recours à une horloge globale logique. Dans un article de 1978, qui est parmi les plus cités, voire le plus cité, de la littérature des systèmes répartis, Lamport définit la relation binaire appelée « précédence causale » (en anglais, happened before relation) entre les évènements1 d’un système, et notée « → », comme étant la plus petite relation satisfaisant les conditions suivantes :

• si a et b sont des évènements sur un même processus, et a arrive avant b, alors a → b ; • si a est l’émission d’un message par un processus et b la réception du même message par un autre

processus, alors a → b ;

• si a → b et b → c, alors a → c.

Cette relation capture la relation de cause à effets.

Ensuite, le temps logique peut être exprimé par un scalaire ou bien par un vecteur, ou encore par une matrice. Les estampilles sont alors des scalaires ou bien des vecteurs, ou encore des matrices. Pour les estampilles contenant un scalaire, leur taille est très petite et fixe, le calcul de l’horloge rapide, et la précision de l’horloge très faible. À l’oppo