EP1649397B1

One-phase commit in a shared-nothing database system

Abstract

This record has no abstract on file.

EP1649397B1, drawing sheet 1
Sheet 1 of 4

Term

Term ended

Expired 28 July 2024, 2.2 years ago.

  1. Priority
  2. Filed
  3. Granted
  4. Expired
  5. Today

12 claims: 12 independent, 0 dependent

  1. 1
    A method for coordinating a distributed transaction in a shared-nothing database system (100), the method comprising:on a first shared-nothing node (110) of said shared-nothing database system (100), a coordinator that is on the first shared-nothing node (110) and that is coordinating the distributed transaction storing, in a first redo log to which the first shared-nothing node (110) writes redo records that pertain to aspects of the distributed transaction that are performed by the first shared-nothing node (110), a commit record that indicates a status of said distributed transaction;wherein the first redo log is stored on a persistent storage device;wherein the persistent storage device is accessible to a participant that is to perform one or more operations as part of said distributed transaction;wherein the participant resides on a second shared-nothing node (150) of said shared-nothing database system (100);wherein the participant writes, to a second redo log that is separate from the first redo log, redo records that pertain to aspects of the distributed transaction that are performed by the second shared-nothing node (150);on the second shared-nothing node (150) of said shared-nothing database system (100), the participant determining the status of said distributed transaction by reading the commit record from the first redo log;andin response to a determination, by the participant, that said status of said distributed transaction is a committed status, the participant committing changes made by the participant as part of the distributed transaction. Un procédé de coordination d'une transaction répartie dans un système de base de données sans partage (100), le procédé comprenant : sur un premier noeud sans partage (110) dudit système de base de données sans partage (100), la conservation en mémoire par un coordinateur qui est sur le premier noeud sans partage (110) et qui coordonne la transaction répartie, dans un premier fichier journal de type redo dans lequel le premier noeud sans partage (110) écrit des enregistrements de type redo qui concernent des aspects de la transaction répartie qui sont exécutés par le premier noeud sans partage (110), d'un enregistrement d'engagement qui indique un état de ladite transaction répartie,dans lequel le premier fichier journal de type redo est conservé en mémoire dans un dispositif à mémoire persistante,dans lequel le dispositif à mémoire persistante est accessible à un participant qui doit exécuter une ou plusieurs opérations en tant que partie prenante de ladite transaction répartie,dans lequel le participant réside sur un deuxième noeud sans partage (150) dudit système de base de données sans partage (100),dans lequel le participant écrit, dans un deuxième fichier journal de type redo qui est distinct du premier fichier journal de type redo, des enregistrements de type redo qui concernent des aspects de la transaction répartie qui sont exécutés par le deuxième noeud sans partage (150),sur le deuxième noeud sans partage (150) dudit système de base de données sans partage (100), la détermination par le participant de l'état de ladite transaction répartie par la lecture de l'enregistrement d'engagement provenant du premier fichier journal de type redo, eten réponse à une détermination, par le participant, que ledit état de ladite transaction répartie est un état engagé, l'engagement par le participant, de modifications effectuées par le participant en tant que partie prenante de la transaction répartie. Verfahren zum Koordinieren einer verteilten Transaktion in einem Shared-Nothing-Datenbanksystem (100) (Datenbanksystem ohne gemeinsame Nutzung), wobei das Verfahren aufweist: auf einem ersten Shared-Nothing-Knoten (110) des Shared-Nothing-Datenbanksystems (100) speichert ein Koordinator, der sich auf dem ersten Shared-Nothing-Knoten (110) befindet und der die verteilte Transaktion koordiniert, in einem ersten Wiederherstellungsprotokoll (Redo-Log), in das der erste Shared-Nothing-Knoten (110) Wiederherstellungseinträge schreibt, die sich auf Aspekte der verteilten Transaktion beziehen, die von dem ersten Shared-Nothing-Knoten (110) ausgeführt werden, einen Festschreibungseintrag, der einen Status der verteilten Transaktion angibt;wobei das erste Wiederherstellungsprotokoll in einer persistenten Speichereinrichtung gespeichert wird;wobei auf die persistente Speichereinrichtung von einem Teilnehmer zugegriffen werden kann, der einen oder mehrere Vorgänge als Teil der verteilten Transaktion ausführen soll;wobei sich der Teilnehmer auf einem zweiten Shared-Nothing-Knoten (150) des Shared-Nothing-Datenbanksystems (100) befindet;wobei der Teilnehmer in ein zweites Wiederherstellungsprotokoll, das von dem ersten Wiederherstellungsprotokoll getrennt ist, Wiederherstellungseinträge schreibt, die sich auf Aspekte der verteilten Transaktion beziehen, die von dem zweiten Shared-Nothing-Knoten (150) ausgeführt werden;auf dem zweiten Shared-Nothing-Knoten (150) des Shared-Nothing-Datenbanksystems (100) bestimmt der Teilnehmer den Status der verteilten Transaktion, indem der Festschreibungseintrag aus dem ersten Wiederherstellungsprotokoll gelesen wird;undansprechend auf eine Feststellung durch den Teilnehmer, dass der Status der verteilten Transaktion ein festgeschriebener Status ist, schreibt der Teilnehmer Änderungen, die von dem Teilnehmer als Teil der verteilten Transaktion durchgeführt worden sind, fest.
  2. 2
    Le procédé selon la Revendication 1 dans lequel :le participant est un premier participant d'une pluralité de participants à ladite transaction répartie,la pluralité de participants comprend un deuxième participant qui n'a pas accès audit dispositif à mémoire persistante,le procédé comprend en outre l'opération d'interaction par le coordinateur avec le deuxième participant selon un protocole d'engagement à deux phases,dans lequel le coordinateur a envoyé un message de préparation au deuxième participant au cours d'une phase de préparation des protocole d'engagement à deux phases du fait que le deuxième participant n'a pas accès audit dispositif à mémoire persistante,dans lequel le coordinateur n'a envoyé aucun message de préparation au premier participant au cours de la phase de préparation du protocole d'engagement à deux phases du fait que le premier participant a accès audit dispositif à mémoire persistante,dans lequel, à un moment où le premier participant est amené à déterminer l'état de ladite transaction répartie par la lecture de l'enregistrement d'engagement provenant du premier fichier journal de type redo, ledit coordinateur a subi une défaillance et n'a pas encore récupéré au moyen d'un processus de récupération quelconque. The method of Claim 1 wherein: the participant is a first participant of a plurality of participants in said distributed transaction;the plurality of participants includes a second participant that does not have access to said persistent storage device;the method further comprises the step of the coordinator interacting with the second participant according to a two-phase commit protocol;wherein the coordinator sent a prepare message to the second participant during a prepare phase of the two-phase commit protocol due to the second participant not having access to said persistent storage device;wherein the coordinator did not send any prepare message to the first participant during the prepare phase of the two-phase commit protocol due to the first participant having access to said persistent storage device;wherein, at a time that the first participant is caused to determine the status of said distributed transaction by reading the commit record from the first redo log, said coordinator has crashed and is not yet being recovered by any recovery process. Verfahren nach Anspruch 1, bei dem: der Teilnehmer ein erster Teilnehmer aus einer Mehrzahl von Teilnehmern an der verteilten Transaktion ist;die Mehrzahl von Teilnehmern einen zweiten Teilnehmer aufweist, der nicht auf die persistente Speichereinrichtung zuzugreifen vermag;das Verfahren ferner den Schritt aufweist, dass der Koordinator mit dem zweiten Teilnehmer gemäß einem Zwei-Phasen-Festschreibungsprotokoll interagiert;wobei der Koordinator während einer Vorbereitungsphase des Zwei-Phasen-Festschreibungsprotokolls eine Vorbereitungsnachricht an den zweiten Teilnehmer gesendet hat, und zwar deshalb, weil der zweite Teilnehmer nicht auf die persistente Speichereinrichtung zuzugreifen vermag;wobei der Koordinator während der Vorbereitungsphase des Zwei-Phasen-Festschreibungsprotokoll keine Vorbereitungsnachricht an den ersten Teilnehmer gesendet hat, und zwar deshalb, weil der erste Teilnehmer auf die persistente Speichereinrichtung zuzugreifen vermag;wobei zu einem Zeitpunkt, an dem der erste Teilnehmer veranlasst wird, den Status den verteilten Transaktion zu bestimmen, indem der Festschreibungseintrag aus dem ersten Wiederherstellungsprotokoll gelesen wird, der Koordinator abgestürzt ist und noch nicht von einem Wiederherstellungsprozess wiederhergestellt wird.
  3. 3
    Le procédé selon la Revendication 1 comprenant en outre les opérations suivantes :l'engagement par le coordinateur de la transaction répartie,après l'engagement par le coordinateur de la transaction répartie, l'envoi par le coordinateur d'un message d'engagement au participant, etl'empêchement d'écraser ou de supprimer l'enregistrement d'engagement jusqu'à ce qu'un ensemble de conditions soit satisfait, une condition dans ledit ensemble de conditions étant que le coordinateur reçoit un message d'accusé de réception d'engagement dudit participant. The method of Claim 1 further comprising the steps of: the coordinator committing the distributed transaction;after the coordinator commits the distributed transaction, the coordinator sending a commit message to the participant;andpreventing the commit record from being overwritten or deleted until a set of conditions is satisfied, wherein one condition in said set of conditions is that the coordinator receives a commit acknowledge message from said participant. Verfahren nach Anspruch 1, ferner die Schritte aufweisend: der Koordinator schreibt die verteilte Transaktion fest,nachdem der Koordinator die verteilte Transaktion festgeschrieben hat, sendet der Koordinator eine Festschreibungsnachricht an den Teilnehmer;undder Festschreibungseintrag wird vor Überschreibung oder Löschung geschützt, bis ein Satz von Bedingungen erfüllt ist, wobei eine Bedingung des Satzes von Bedingungen ist, dass der Koordinator eine Festschreibungs-Bestätigungsnachricht von dem Teilnehmer erhält.
  4. 4
    Le procédé selon la Revendication 1 comprenant en outre les opérations suivantes :l'envoi par le participant d'une première information au coordinateur, la première information étant associée à une tâche exécutée par ledit participant en tant que partie prenante de ladite transaction répartie, etl'exécution par le coordinateur d'une comparaison entre la première information et des informations associées au deuxième fichier journal de type redo, etla détermination par le coordinateur s'il convient d'engager la transaction en fonction, au moins en partie, de ladite comparaison. The method of Claim 1 further comprising the steps of: the participant sending a first piece of information to the coordinator, wherein the first piece of information is associated with work performed by said participant as part of said distributed transaction;andthe coordinator performing a comparison between the first piece of information and information associated with the second redo log;andthe coordinator determining whether to commit the transaction based, at least in part, on said comparison. Verfahren nach Anspruch 1, ferner die Schritte umfassend: der Teilnehmer sendet erste Informationselemente an den Koordinator, wobei die ersten Informationselemente Vorgängen zugeordnet sind, die als Teil der verteilten Transaktion von dem Teilnehmer ausgeführt werden;undder Koordinator führt einen Vergleich zwischen den ersten Informationselementen und Informationen, die dem zweiten Wiederherstellungsprotokoll zugeordnet sind, aus;undder Koordinator bestimmt, zumindest zum Teil beruhend auf dem Vergleich, ob die Transaktion festgeschrieben werden soll.
  5. 5
    Le procédé selon la Revendication 4 dans lequel l'information contient un numéro de séquence de fichier journal de la modification la plus récente effectuée par le participant en tant que partie prenante de la transaction répartie. The method of Claim 4 wherein the piece of information includes a log-sequence-number of the latest change made by the participant as part of the distributed transaction. Verfahren nach Anspruch 4, bei dem die Informationselemente eine Protokoll-Sequenznummer der letzten Änderung aufweisen, die von dem Teilnehmer als Teil der verteilten Transaktion ausgeführt worden ist.
  6. 6
    Le procédé selon la Revendication 5 dans lequel l'opération d'envoi comprend les opérations suivantes :l'identification par le participant d'un message qui est envoyé audit premier noeud sans partage (110) dans un but sans relation avec la transaction répartie, etl'insertion du numéro de séquence de fichier journal dans ledit message. The method of Claim 5 wherein the step of sending includes the steps of: the participant identifying a message that is being sent to said first shared-nothing node (110) for a purpose unrelated to the distributed transaction;andpiggybacking the log-sequence number on said message. Verfahren nach Anspruch 5, bei dem der Schritt des Sendens die Schritte aufweist: der Teilnehmer erkennt eine Nachricht, die an den ersten Shared-Nothing-Knoten (110) zu einem Zweck, der nicht mit der verteilten Transaktion in Beziehung steht, gesendet wird;unddie Protokoll-Sequenznummer wird Huckepack zu der Nachricht hinzugefügt.
  7. 7
    A method for performing a distributed transaction in a shared-nothing database system (100), the method comprising:assigning a participant to perform one or more operations as part of said distributed transaction;wherein the participant resides on a second shared-nothing node (150) of said shared-nothing system (100);said participant storing in a persistent storage device, in a redo log to which only the participant writes redo records that pertain to aspects of the distributed transaction that are performed by the participant, status information that indicates changes made to data in the shared-nothing database on which said one or more operations operate, by the participant during performance of said one or more operations;wherein the persistent storage device is accessible to a coordinator that is responsible for coordinating said distributed transaction;wherein the coordinator resides on a first shared-nothing node (110) of said shared-nothing database system (100);on said first shared-nothing node (110) of said shared-nothing database system (100), said coordinator determining, based on the status information contained in the redo log on said persistent storage device, whether the participant has written to persistent storage changes produced by performance of the one or more operations;andthe coordinator determining whether the distributed transaction can be committed based, at least in part, on whether the participant has written to persistent storage changes produced by performance of the one or more operations. Un procédé d'exécution d'une transaction répartie dans un système de base de données sans partage (100), le procédé comprenant : la nomination d'un participant destiné à exécuter une ou plusieurs opérations en tant que partie prenante de ladite transaction répartie,dans lequel le participant réside sur un deuxième noeud sans partage (150) dudit système sans partage (100),la conservation en mémoire par ledit participant, dans un dispositif à mémoire persistante, dans un fichier journal de type redo dans lequel uniquement le participant écrit des enregistrements de type redo qui concernent des aspects de la transaction répartie qui sont exécutés par le participant, d'informations d'état qui indiquent des modifications apportées à des données de la base de données sans partage sur lesquelles lesdites une ou plusieurs opérations agissent, par le participant au cours de l'exécution desdites une ou plusieurs opérations,dans lequel le dispositif à mémoire persistante est accessible à un coordinateur qui est chargé de la coordination de ladite transaction répartie,dans lequel le coordinateur réside sur un premier noeud sans partage (110) dudit système de base de données sans partage (100),sur ledit premier noeud sans partage (110) dudit système de base de données sans partage (100), la détermination par ledit coordinateur, en fonction des informations d'état contenues dans le fichier journal de type redo sur ledit dispositif à mémoire persistante, si le participant a écrit vers une mémoire persistante des modifications produites par l'exécution des une ou plusieurs opérations, etla détermination par le coordinateur si la transaction répartie peut être engagée en fonction, au moins en partie, du fait que le participant a écrit vers une mémoire persistante des modifications produites par l'exécution des une ou plusieurs opérations. Verfahren zum Ausführen einer verteilten Transaktion in einem Shared-Nothing-Datenbanksystem (100) (Datenbanksystem ohne gemeinsame Nutzung), wobei das Verfahren aufweist: Beauftragen eines Teilnehmers mit der Ausführung eines oder mehrerer Vorgänge als Teil der verteilten Transaktion;wobei der Teilnehmer sich auf einem zweiten Shared-Nothing-Knoten (150) des Shared-Nothing-Systems (100) befindet;der Teilnehmer speichert, in einer persistenten Speichereinrichtung, Statusinformationen in einem Wiederherstellungsprotokoll (Redo-Log), in das nur der Teilnehmer Wiederherstellungseinträge einschreibt, die sich auf Aspekte der verteilten Transaktion beziehen, die von dem Teilnehmer ausgeführt werden, wobei die Statusinformationen Änderungen angeben, die an Daten in der Shared-Nothing-Datenbank, auf die der eine oder die mehreren Vorgänge einwirken, vorgenommen werden, und zwar durch den Teilnehmer während der Ausführung des einen oder der mehreren Vorgänge;wobei ein Koordinator, der für die Koordinierung der verteilten Transaktion verantwortlich ist, auf die persistente Speichereinrichtung zuzugreifen vermag;wobei der Koordinator sich auf einem ersten Shared-Nothing-Knoten (110) des Shared-Nothing-Datenbanksystems (100) befindet;auf dem ersten Shared-Nothing-Knoten (110) des Shared-Nothing-Datenbanksystems (100) bestimmt der Koordinator, beruhend auf den Statusinformationen, die in dem Wiederherstellungsprotokoll in der persistenten Speichereinrichtung enthalten sind, ob der Teilnehmer Änderungen, die durch die Ausführung des einen oder der mehreren Vorgänge hervorgerufen werden, in persistenten Speicher eingeschrieben hat;undder Koordinator bestimmt, ob die verteilte Transaktion festgeschrieben werden kann, und zwar zumindest zum Teil beruhend darauf, ob der Teilnehmer Änderungen, die durch die Ausführung des einen oder der mehreren Vorgänge hervorgerufen werden, in persistenten Speicher eingeschrieben hat.
  8. 8
    Le procédé selon la Revendication 7 dans lequel :l'opération de conservation en mémoire par ledit participant, sur un dispositif à mémoire persistante, d'informations d'état qui indiquent des modifications effectuées par le participant au cours de l'exécution desdites une ou plusieurs opérations comprend la conservation en mémoire par ledit participant d'informations de type redo dans un fichier journal de type redo sur ledit dispositif à mémoire persistante, etl'opération de détermination par ledit coordinateur, en fonction des informations d'état sur ledit dispositif à mémoire persistante, si le participant a écrit vers une mémoire persistante des modifications produites par l'exécution des une ou plusieurs opérations comprend une inspection du fichier journal de type redo du participant de façon à déterminer si les informations de type redo pour lesdites modifications ont été écrites dans ladite mémoire persistante. The method of Claim 7 wherein: the step of said participant storing, on a persistent storage device, status information that indicates changes made by the participant during performance of said one or more operations includes said participant storing redo information in a redo log on said persistent storage device;andthe step of said coordinator determining, based on the status information on said persistent storage device, whether the participant has written to persistent storage changes produced by performance of the one or more operations includes inspecting the redo log of the participant to determine whether the redo information for said changes have been written to said persistent storage. Verfahren nach Anspruch 7, bei dem: der Schritt, dass der Teilnehmer in einer persistenten Speichereinrichtung Statusinformationen speichert, die Änderungen angeben, die von dem Teilnehmer während der Ausführung des einen oder der mehreren Vorgänge vorgenommen werden, aufweist, dass der Teilnehmer Wiederherstellungsinformationen in einem Wiederherstellungsprotokoll in der persistenten Speichereinrichtung speichert;undder Schritt, dass der Koordinator beruhend auf den Statusinformationen in der persistenten Speichereinrichtung bestimmt, ob der Teilnehmer Änderungen, die durch die Ausführung des einen oder der mehreren Vorgänge hervorgerufen werden, in persistenten Speicher eingeschrieben hat, umfasst, das Wiederherstellungsprotokoll des Teilnehmers zu sichten, um zu bestimmen, ob die Wiederherstellungsinformationen für die Änderungen in den persistenten Speicher eingeschrieben worden sind.
  9. 9
    Le procédé selon la Revendication 7 dans lequel :le participant est un premier participant d'une pluralité de participants à ladite transaction répartie,la pluralité de participants comprend un deuxième participant qui conserve en mémoire des informations d'état sur un deuxième dispositif à mémoire persistante qui n'est pas accessible par ledit coordinateur, etle procédé comprend en outre l'opération d'interaction par le coordinateur avec le deuxième participant selon un protocole d'engagement à deux phases. The method of Claim 7 wherein: the participant is a first participant of a plurality of participants in said distributed transaction;the plurality of participants includes a second participant that stores status information on a second persistent storage device that is not accessible by said coordinator;andthe method further comprises the step of the coordinator interacting with the second participant according to a two-phase commit protocol. Verfahren nach Anspruch 7, bei dem: der Teilnehmer ein erster Teilnehmer aus einer Mehrzahl von Teilnehmern an der verteilten Transaktion ist;die Mehrzahl von Teilnehmern einen zweiten Teilnehmer aufweist, der Statusinformationen in einer zweiten persistenten Speichereinrichtung speichert, auf die der Koordinator nicht zuzugreifen vermag;unddas Verfahren ferner den Schritt aufweist, dass der Koordinator mit dem zweiten Teilnehmer gemäß einem Zwei-Phasen-Festschreibungsprotokoll interagiert.
  10. 10
    Le procédé selon la Revendication 7 dans lequel, les informations sur ledit dispositif à mémoire persistante indiquent que le participant n'a pas écrit vers une mémoire persistante des modifications produites par l'exécution des une ou plusieurs opérations, et le procédé comprend en outre l'envoi par le coordinateur d'un message de type redo forcé au participant de façon à amener le participant à écrire vers une mémoire persistante les modifications produites par l'exécution des une ou plusieurs opérations. The method of Claim 7 wherein:the information on said persistent storage device indicates that the participant has not written to persistent storage changes produced by performance of the one or more operations;andthe method further comprises the coordinator sending a force redo message to the participant to cause the participant to write to persistent storage the changes produced by performance of the one or more operations. Verfahren nach Anspruch 7, bei dem: die Informationen in der persistenten Speichereinrichtung angeben, dass der Teilnehmer keine Änderungen, die durch die Ausführung des einen oder der mehreren Vorgänge hervorgerufen werden, in persistenten Speicher eingeschrieben hat;unddas Verfahren ferner aufweist, dass der Koordinator eine Nachricht zur erzwungenen Wiederherstellung an den Teilnehmer sendet, um den Teilnehmer zu veranlassen, die Änderungen, die durch die Ausführung des einen oder der mehreren Vorgänge hervorgerufen werden, in persistenten Speicher einzuschreiben.
  11. 11
    Le procédé selon la Revendication 10 dans lequel l'opération d'envoi d'un message de type redo forcé comprend les opérations suivantes :l'identification d'un message qui est envoyé audit deuxième noeud sans partage (150) dans un but sans relation avec la transaction répartie, etl'insertion du message de type redo forcé dans ledit message. The method of Claim 10 wherein the step of sending a force redo message includes the steps of: identifying a message that is being sent to said second shared-nothing node (150) for a purpose unrelated to the distributed transaction;andpiggybacking the force redo message on said message. Verfahren nach Anspruch 10, bei dem der Schritt des Sendens einer Nachricht zur erzwungenen Wiederherstellung die Schritte aufweist: Erkennen einer Nachricht, die an den zweiten Shared-Nothing-Knoten (150) zu einem Zweck, der nicht mit der verteilten Transaktion in Beziehung steht, gesendet wird;unddie Nachricht zum erzwungenen Wiederherstellen wird Huckepack zu der Nachricht hinzugefügt.
  12. 12
    A computer-readable medium carrying one or more sequences of instructions which, when executed by one or more processors (404), causes the one or more processors (404) to perform a method as recited in any of Claims 1-11. Computerlesbares Medium, dass eine oder mehrere Befehlssequenzen aufweist, die, wenn sie von einem oder mehreren Prozessoren (404) ausgeführt wird/werden, den einen oder die mehreren Prozessoren (404) dazu veranlasst/veranlassen, das Verfahren nach einem der Ansprüche 1-11 auszuführen. Un support lisible par ordinateur contenant une ou plusieurs séquences d'instructions qui, lorsqu'elles sont exécutées par un ou plusieurs processeurs (404), amènent les un ou plusieurs processeurs (404) à exécuter un procédé selon l'une quelconque des Revendications 1 à 11.