Distributed data storage
2 claims: 1 independent, 1 dependent
- 1PATENTKRAV 1. Metod för underhåll av data i ett datalagringssystem, vilket innefattar ett flertal datalagringsnoder, varvid metoden används i en lagringsnod i 5 datalagringssystemet och innefattar:-övervakning av statusen (59) hos andra lagringsnoder i systemet såväl som skrivoperationer (65, 67, 69) som utförs i datalagringssystemet;-detektering (51), på basis av övervakandet, av tillstånd i datalagringssystemet som implicerar behovet av replikering av data mellan 10 noderna i datalagringssystemet;och -initiering av en replikeringsprocess (53) i det fall att ett sådant tillstånd detekteras, varvid replikeringsprocessen inkluderar sändning av ett IPmulticast-meddelande till ett flertal lagringsnoder, vilket meddelande frågar vilka av dessa lagringsnoder som lagrar särskilda data. 15 2. Metod enligt krav 1, varvid övervakandet innefattar lyssnandet (59) till heartbeat-signaler från andra lagringsnoder i systemet, och varvid ett tillstånd som implicerar behovet av replikering är en fallerande lagringsnod. 3. Metod enligt krav 1 eller 2, varvid nämnda data inkluderar filer och ett tillstånd är endera av en filborttagning och en filkonflikt. 20 4. Metod enligt något av föregående krav, varvid en repiikeringslista, som innefattar filer i behov av replikering, underhålles. 5. Metod enligt krav 4, varvid replikeringslistan inkluderar prioriteter. 6. Metod enligt något av föregående krav, varvid replikeringsprocessen vidare inkluderar: 25 -mottagning av svar från de lagringsnoder som innehåller nämnda särskilda data;-bestämning av huruvida nämnda särskilda data lagras på ett tillräckligt antal lagringsnoder;och -om så inte är fallet, val av åtminstone en ytterligare lagringsnod och 30 sändning av nämnda särskilda data till denna lagringsnod. 7. Metod enligt krav 6, vidare innefattande uppdatering av nämnda särskilda data på lagringsnoder som innehåller obsoleta versioner därav. 533 007 8. Metod enligt krav 6 eller 7, varvid replikeringsprocessen inleds med att lagringsnoden försöker bli master för filen som ska replikeras, bland alla lagringsnoder i systemet. 9. Metod enligt något av föregående krav, varvid övervakandet vidare inkluderar övervakning av läsoperationer (63) som utförs i datalagringssystemet. 10. Datalagringsnod för underhåll av data i ett datalagringssystem, som innefattar ett flertal datalagringsnoder, vilken datalagringsnod innefattar: -medel för övervakning av statusen hos andra lagringsnoder i systemet såväl som skrivoperationer som utförs i datalagringssystemet;- medel för detektering, på basis av övervakandet, tillstånd i datalagringssystemet, vilka implicerar behovet av replikering av data mellan noderna i datalagringssystemet;och -medel för initiering av en replikeringsprocess ifall ett sådant tillstånd detekteras, varvid replikeringsprocessen inkluderar sändning av ett IPmulticast-meddelande till ett flertal lagringsnoder, vilket meddelande frågar vilka av dessa lagringsnoder som lagrar särskilda data. 11. Metod för skrivning av data till ett datalagringssystem, vilket innefattar ett flertal datalagringsnoder, vilken metod används i en server, som kör en applikation, vilken åtkommer data i datalagringssystemet, och innefattar: -sändning (41) av en IP-multicast-lagringsförfrågan till ett flertal bland lagringsnoderna;-mottagning (43) av ett flertal svar från en delmängd bland nämnda lagringsnoder, varvid svaren inkluderar geografiska data avseende varje lagringsnods geografiska position;-val (45) av åtminstone två lagringsnoder i delmängden, på basis av nämnda svar;och -sändning (47) av data och en dataidentifierare, motsvarande dessa data, till de valda lagringsnoderna. 12. Metod enligt krav 11, varvid den geografiska positionen inkluderar latituden och longituden för lagringsnoden i fråga. 533 007 13. Metod enligt krav 12, varvid svaren vidare inkluderar systemåldern för lagringsnoden i fråga. 14. Metod enligt krav 12 eller 13, varvid svaren vidare inkluderar systemlasten för lagringsnoden i fråga. 15. Metod enligt något av krav 12-14, varvid nämnda IP-multicastlagringsförfrågan inkluderar en dataidentifierare, som identifierar nämnda data som ska lagras. 16. Metod enligt något av krav 12-15, varvid åtminstone tre noder väljs. 17. Metod enligt något av krav 12-16, varvid en lista med lagringsnoder, vilka framgångsrikt lagrat nämnda data, skickas till de valda lagringsnoderna. 18. Server, anordnad för skrivning av data till ett datalagringssystem och innefattande ett flertal datalagringsnoder, vilken server innefattar: -medel för sändning av en IP-multicast-lagringsförfrågan till ett flertal bland nämnda lagringsnoder;-medel för mottagning av ett flertal svar från en delmängd bland nämnda lagringsnoder, varvid svaren inkluderar geografiska data avseende varje lag ringsnods geografiska position;-medel för val av åtminstone två lagringsnoder i delmängden på basis av nämnda svar;och -medel för sändning av data och en dataidentifierare, som motsvarar dessa data, till den valda lagringsnoden. 533 007 533 007
- 22/7
Independent claims2
116 paragraphs in 1 section, as filed
(12) Patent No. SE 533 007 C2 rt<sup>6LS</sup>
Μ ·· o
<img file="SE533007C2_D0001.tif" />
Sweden (21) Patent Application Number. 0802277-4 (45) Patent granted: 2010-06-08 (41) Application generally available: 2010-04-25 (22) Patent application submitted: 2008-10-24 (24) Maturity date: 2008-10-24 (83) Deposit of microorganism: - (30) Priority information: - (51) International class:
G06F15 / 173 (2006.01) (73) Patent holder: ILT Productions AB, Box 177, 371 22 Karlskrona SE
<td>(72) Inventor:</td><td>Christian Melander, Rödeby SE Stefan Bembo, Karlskrona SE Gustav Petersson, Karlskrona SE Roger Persson, Karlskrona SE</td>
<td>(74) Agents:</td><td>AWAPATENT AB, Box 5117, 200 71 Malmö SE</td>
<td>(54) Name:</td><td>Distributed data storage</td>
<td>(56) Quoted publications:</td><td>US 20040059805 Al</td>
<td>(47) Summary:</td><td>The present invention relates to a distributed data storage system comprising a plurality of storage nodes. When using unicast and multicast transmission, a server application can read and write data in the storage system. Each storage node can monitor the read and write operations of the system as well as the status of other storage nodes. In this way, the storage nodes can detect a need for replication of files in the system, and can perform a replication process that serves to maintain storage of a sufficient number of copies of files with correct versions at geographically separated locations.</td>
<img file="SE533007C2_D0002.tif" />
533 007
SUMMARY
The present invention relates to a distributed data storage system comprising a plurality of storage nodes. When using unicast and multicast transmission, a server application can read and write data in the storage system. Each storage node can monitor the read and write operations of the system as well as the status of other storage nodes. In this way, the storage nodes can detect a need for replication of files in the system, and can perform a replication process that serves to maintain storage of a sufficient number of copies of files with correct versions at geographically separated locations.
533 007
Technical area
The present disclosure relates to methods for writing and maintaining data in a data storage system which includes a plurality of data storage nodes, the methods being used in a server and in a storage node in the data storage system. The description further relates to storage nodes or servers, which are capable of performing such methods.
Background
Such a method is disclosed, for example, in US, 2005/0246393, A1. This method is shown for a system that uses multiple storage centers in geographically separated locations. Distributed object storage controllers are included to maintain stored data information.
One problem associated with such a system is how it achieves simple yet robust data writing and maintenance.
Summary of the Invention
It is therefore an object of the present description to provide robust writing or maintenance of data in a distributed storage system, without the use of centralized maintenance servers, which themselves can be a weak link in a system. This object is achieved by a method of the initially mentioned type, which is provided in a storage node and comprises: monitoring the status of other storage nodes in the system as well as writing operations performed in the data storage system, detecting, on the basis of monitoring, conditions in the data storage system which implicate the need for replication of data between the nodes in the data storage system, and initiation of a replication process in case such a condition is detected. The replication process includes sending an IP multicast message to a plurality of storage nodes including a request for which of these storage nodes stores particular data.
533 007
Using such a method, each storage node can be active in maintaining data throughout the system. In the event that a storage node fails, its data can be recreated by other nodes in the system, so the system can be considered self-healing.
Monitoring may include listening to heartbeat signals from other storage nodes in the system. A condition that implies the need for replication may then be a failing storage node.
Said data includes files, and a state that implies the need for replications may then be either a file deletion or a file conflict.
A replication list, which includes files in need of replication, can be maintained and can include priorities.
The replication process may further include: receiving responses from the storage nodes containing said particular data, determining whether said particular data is stored in a sufficient number of storage nodes, and, if not, selecting at least one additional storage node and transmitting said particular data to this storage node. Furthermore, said particular data on storage nodes, which contain obsolete versions thereof, can be updated.
Furthermore, the replication process can begin with the storage node trying to become the master, among all the storage nodes in the system, for the file to be replicated.
The monitoring may further include monitoring of read operations performed in the data storage system.
The present description further relates to a data storage node for performing data maintenance which corresponds to the method. The storage node then generally comprises means for performing the steps of the method.
The object is also achieved by a method for writing data to a data storage system of the type mentioned in the introduction, which method is provided in a server running an application that accesses data in the data storage system. The method comprises: transmitting an IP multicast storage request to a plurality of storage nodes, receiving a plurality of responses from a subset of said storage nodes, the replies including geographical data regarding the geographical position of each storage node, selection of at least two storage nodes in the subset, based on said answer, as well
533 007 transmits data and a data identifier corresponding to said data to the selected storage nodes.
This method achieves robust data writing by efficiently achieving geographical diversity.
The geographical position may include latitude and longitude of the storage node in question, and the responses may further include system load and / or system age of the storage node in question.
Said IP multicast storage request may include a data identifier, which identifies the data to be stored.
Typically, at least three nodes can be selected for storage, and a list of storage nodes that successfully store said data can be sent to the selected storage nodes.
The present description further relates to a server for performing data writing which corresponds to the method. The server then generally comprises means for performing the steps of the method.
Brief description of the details
Fig. 1 illustrates a distributed data storage system.
Figs. 2A -2C, and Fig. 3 illustrate a data reading process.
Figs. 4A -4C, and Fig. 5 illustrate a data writing process.
Fig. 6 schematically illustrates a situation where a number of files are stored at a number of data storage nodes.
Fig. 7 illustrates the transmission of heartbeat signals.
Fig. 8 is an overview of a data maintenance process.
Detailed description
The present disclosure relates to a distributed data storage system comprising a plurality of storage nodes. The structure of the system and the context in which it is used are shown in Figure 1.
Thus, a user's computer 1 accesses, via the Internet 3, an application 5 running on a server 7. The user context, as shown here, is thus a common client-server configuration, which is well known in itself. It should
533 007, however, it is noted that the data storage system that will be displayed is also useful in other configurations.
In the illustrated case, two applications 5, 9 are run on the server 7. Of course, this number of applications may be different. Each application has an Application Programming Interface (API) 11 that provides an interface to the distributed data storage system 13 and supports requests (typically write and read requests) from applications running on the server. From the application's point of view, reading or writing information from / to the data storage system 13 need not appear as different from the use of any other type of storage solution, for example a file server or simply a hard disk.
Each AP111 communicates with storage nodes 15 in the data storage system 13, and the storage nodes communicate with each other. These communications are based on TCP (Transmission Control Protocol) and UDP (User Datagram Protocol). These concepts are well known to those skilled in the art and need not be explained further.
It should be noted that different APIs 11 on the same server 7 can access different sets of storage nodes 15. It should further be noted that there may be more than one server 7 accessing each storage node 15. However, this does not affect the way the storage nodes work , as will be described later.
The components of the distributed data storage system are the storage nodes 15 and the APLs 11 of the server 7 accessing the storage nodes
15th The present description therefore relates to methods performed in the server 7 and in the storage nodes 15. These methods will primarily be implemented as software, which are run on the server and the storage nodes respectively, which together determine how the data storage system will work and what features it has.
The storage node 15 can typically be executed as a file server, which can typically be provided with a number of functional blocks. Thus, the storage node may comprise a storage medium 17, which typically includes a plurality of hard drives, optionally configured as a RAID (Redundant Array of
533 007
Independent Disk) system. However, other types of storage media are also conceivable.
The storage node 15 may further comprise a register 19 which includes lists of data entity / storage node relationships such as a host list (host list), which will be described later.
In addition to the host list, each storage node further includes a node list, which includes the IP addresses of all storage nodes in its set or group of storage nodes. The number of storage nodes in a group can vary from a few to hundreds of storage nodes. The node list may further have a version number.
Further, the storage node 15 may include a replication block 21 and a cluster monitoring block 23. The replication block 21 includes a storage node API 25, and is configured to execute functions for identifying the need for, and performing a replication process, as will be described in detail later. The storage node API 25 in the replication block 21 may include code that extends substantially to the code of the server node's storage node API 11, since the replication process includes steps that largely correspond to the steps performed by server 7 during the read and write operations that occur. to be described. For example, the write operation performed during replication largely corresponds to the write operation performed by the server 7. Cluster monitoring block 23 is configured to perform monitoring of other storage nodes in the data storage system 13, as will be described in more detail later.
The distributed data storage system storage nodes 15 can be considered to be at the same hierarchical level. There is no need to designate any master storage node, which is responsible for maintaining a register of stored data entities and for monitoring data conflicts, etc. Instead, all storage nodes may be considered equal and may occasionally perform data management operations for other storage nodes in the system. This equality ensures that the system is robust. In the event that a storage node fails, other nodes in the system will cover the failing node and ensure reliable data storage.
The operation of the system will be described in the following order: reading data, writing data and data maintenance. Although these methods work
533 007 well together, it should be noted that they can in principle be performed independently of one another. For example, the data storage method provides excellent properties even if the data writing method is used and vice versa.
The reading method is now described with reference to Figures 2A-2C and 3, the latter of which is a flow chart illustrating the method.
The reading, as well as other features of the system, uses multicast communication to communicate simultaneously with a plurality of storage nodes. By multicast or IP multicast is meant here point-to-multipoint communication which is accomplished by sending a message to an IP address reserved for multicast applications.
For example, a message, typically a request, is sent to such an IP address (for example, 244.0.0.1), and a number of receiving servers are registered as subscribers for that IP address. Each of the recipients has their own IP address. When a switch on the network receives the message addressed to 244.0.0.1, the switch forwards the message to the IP addresses of each server registered as a subscriber.
In principle, only one server can be registered as a subscriber for a multicast address, in which case point-to-point communication is achieved. However, in this context such a communication is nevertheless considered a multicast communication, since a multicast method is used.
Unicast communication is also used and refers to a communication with a single receiver.
Referring to Figs. 2A and 3, the method for retrieving data from a data storage system comprises transmitting 31 a multicast request to a plurality of storage nodes 15.1. In the illustrated case, there are five storage nodes each having an IP (Internet Protocol) address 192.168. 1.1,192.168.1.2, etc. Of course, the number of storage nodes is just one example. The request includes a data identifier "2B9B4A97-76E5-499E-A21A6D7932DD7927", which may for example be a Universally Unique Identifier (UUID), which is well known in itself.
The storage nodes scan themselves for data that matches the identifier. If such data is found, a storage node sends a response, which is received 33 by server 7, cf. Fig. 2B. As shown
533 007, the answer may optionally contain additional information in addition to an indication that the storage node has a copy of relevant data. Specifically, the response may contain information from the storage node's register regarding other storage nodes containing said data, information regarding which version of said data contained in the storage node, and information about which load the storage node is currently subjected to.
On the basis of the responses, the server selects one or more storage nodes from which data is to be retrieved and 37 sends a unicast request for data to this storage node (s), cf. Fig. 2C.
In response to the request for data, the storage node (s) sends relevant data using unicast to the server, which receives 39 said data. In the illustrated case, only one storage node is selected. While this is sufficient, it is possible to select more than one storage node for the purpose of receiving two sets of data, allowing for a consistency check. If the data transfer fails, the server may select another retrieval node for retrieval.
The choice of storage nodes can be based on an algorithm that takes into account several factors in order to obtain good overall system performance. Typically, the storage node that has the latest version and the lowest load is selected, although other concepts are fully conceivable.
Optionally, the operation can be completed by sending the server a list to all the storage nodes involved, which indicates which nodes contain said data and which version. On the basis of this information, the storage nodes themselves can maintain data with the replication process that will be described.
Figures 4A-4C and Figure 5 illustrate a data writing process for the distributed data storage system.
Referring to Figs. 4A and 5, the method comprises a server transmitting 41 a multicast storage request to a plurality of storage nodes. The storage request includes an identifier and basically consists of a query as to whether receiving storage nodes can store this file. Optionally, the storage nodes can check in their internal registers whether they already have a file
533 007 with this name, and can notify server 7 in the unlikely event that this is the case, so that the server can rename the file.
In any case, at least a subset of the storage nodes will provide responses via unicast transmissions to the server 7. Typically, storage nodes with a predetermined minimum free disk space respond to the request. The server 7 receives 43 responses, which include geographic data for each server's geographic location. For example, as indicated in Fig. 4B, said geographical data may include latitude, longitude and altitude for each server. Other types of geographical data may also be conceivable, such as postal code or the like.
In addition to said geographical data, additional information may be provided, which serves as input to a storage node selection process. In the example shown, the amount of free storage space is provided in each storage node, along with an indication of the storage node's system age and the load to which the storage node is currently exposed.
On the basis of the responses received, the server selects at least two, in a typical embodiment three, storage nodes in the subset, for storing data. The selection of storage nodes is performed using an algorithm that takes into account different data. The selection is carried out in order to achieve some kind of geographical diversity. At the very least, it should be avoided that only file servers in the same rack are selected as storage nodes. Typically, extensive geographical diversity can also be achieved including selection of storage nodes on different continents. In addition to geographical diversity, other parameters can be included in the selection algorithm. As long as some minimal geographical diversity is achieved, free storage space, system age and current load can also be taken into account.
When the storage node is selected, the data to be stored and a corresponding data identifier are sent to each node, typically with the help of a unicast transmission.
Optionally, the operation can be terminated by each storage node that has successfully performed the write operation sends a acknowledgment signal to the server. The server then sends a list to all storage nodes involved, which indicates which nodes have successfully written said data and
533 007 who haven't done it. On the basis of this information, the storage nodes themselves can properly maintain said data by means of the replication process which will be described. For example, if a storage node's writing failed, there is a need to replicate the file to an additional storage node in order to achieve the desired number of storage nodes for that file.
The data writing method, in itself, allows an API in a server 7 to store data in a very robust way, as excellent geographical diversity can be achieved.
In addition to the write and read operations, the API in server 7 can perform operations that delete files and update files. These processes will be described in conjunction with the data maintenance process described below.
The purpose of the data maintenance process is to ensure that a reasonable number of non-failing storage nodes each store the latest version of each file. Furthermore, it can provide the feature that no deleted files are stored at any storage node. Maintenance is performed by the storage nodes themselves. Thus, there is no need for any designated "master" who takes responsibility for the maintenance of the data storage. This ensures improved reliability, as the master itself would otherwise be a weak point in the system.
Fig. 6 schematically illustrates a situation where a number of files are stored among a number of data storage nodes. In the illustrated case, as an example, twelve nodes are shown which have continuously numbered IP addresses from 192.168.1.1 to 192.168.1.12. Of course, however, the IP addresses do not have to be in the same range at all. Each node stores one or two files, which for simplicity are identified by the letters AF.
Referring to Fig. 8, the method of maintaining the data comprises the detection 51 of states in the data storage system, which implies the need for replication of data between the nodes in the data storage system, and a replication process 53. The result of the detection process 51 is a list 55 of files for which the replication has been identified. The list may further include data regarding the order of priority for different replication needs. On the basis of this list, the replication process 53 is performed.
533 007
The robustness of the distributed storage is based on the fact that a reasonable number of copies of each file, correct versions, are stored in the system. In the illustrated case, three copies of each file are stored. However, if, for example, the storage node with the address 192.168.1.5 fails, the desired number of stored files for “B” and “C would be underestimated.
An event that results in the need for replication is because a storage node in the system stops working.
Each storage node in the system can monitor the status of other storage nodes in the system. This can be done by allowing each storage node to output a so-called heartbeat signal at regular intervals, as shown in Fig. 7.1. The illustrated case, the storage node with address 192.168.1.7 outputs a multicast signal 57 to other storage nodes in the system, indicating correctly. This signal can be received by all other functioning storage nodes in the system performing heartbeat monitoring 59 (cf. Figure 8), or a portion thereof. However, in the case of the storage node with address 192,168.1.5, this node fails and it does not emit a heartbeat signal. Thus, the other storage nodes will detect that no heartbeat signal has been emitted by this node for a long time, indicating that the storage node in question is disconnected.
In addition to the storage node's address, the heartbeat signal may include its node list version number. Another storage node, which listens to the heartbeat signal and finds that the sending storage node has a later version of the node list, can then request that the sending storage node transmit its node list. This means that adding or deleting storage nodes can be accomplished simply by adding or removing a storage node and sending a new node list version to a single storage node. This storage node will then spread to all other storage nodes in the system.
Again with reference to Fig. 8, 61 searches each storage node in its internal register for files stored in the failing storage node. Storage nodes, which themselves store files "B" and "C", will find the failing storage node and can therefore add the corresponding file to their lists 55.
533 007
However, the detection process may also detect other conditions that imply the need for replication of a file. Typically, such a condition can be a conflict, that is, one or more storage nodes have an obsolete version of the file. A removal operation also gives it a replication process, as this process can perform the effective physical removal of the file. The server's deletion operation then only needs to ensure that the storage nodes set a deletion flag for the file in question. Each node can therefore monitor read and write operations performed in the data storage system. Information provided by server 7 at the end of read-responsive write operations may indicate that a storage node contains an obsolete version of a file (in the case of a read operation) or that a storage node has not successfully completed a write operation. In both cases, there is a need for data maintenance through replication so that the overall objectives of the maintenance process are met.
In addition to the basic read and write operations 63, 65, at least two additional processes may provide indications of a need for replication, namely the removal 67 and the update processes, which will now be briefly described.
The removal process is initiated by the server 7 (cf. Fig. 1). As with the read process, the server sends a request using multicasting to all storage nodes, in order to find out which storage nodes have data with a particular data identifier. The storage nodes scan themselves for data with the relevant identifier, and respond with a unicast broadcast if they have the data in question. The answer may include a list, from the storage node's register, of other storage nodes having said data. Server 7 then sends a unicast request, to the storage nodes that are supposed to store the file, that the file should be deleted. Each storage node sets a flag with respect to the file and indicates that it should be deleted. The file is then added to the replication list, and a confirmation is sent to the server. The replication process then physically deletes the file, as will be described.
The update process has a search function, similar to the deletion process, and a write function, similar to the one performed in the writing process. The server sends a request using multicasting to all storage nodes, i
533 007 purpose of finding out which storage nodes have data with a particular data identifier. The storage nodes scan themselves for data with the relevant identifier, and respond with a unicast broadcast if they have the data in question. The answer may include a list, from the storage node's register, of other storage nodes containing said data. The server 7 then sends a unicast request, which orders the storage nodes to update said data. The request, of course, contains the updated data. The storage nodes update said data and send a receipt to the server, which responds by sending a unicast broadcast containing a list of the storage nodes that have successfully updated said data, and storage nodes that did not. Again, this list can be used by the maintenance process.
Again with reference to Fig. 8, read 63, write 65, removal 67, and update operations 69 can all indicate that a need for replication exists. The same is true for heartbeat monitoring 59. Thus, the overall detection process 51 generates data regarding which files need replication. For example, a read or update operation may show that a particular storage node contains an obsolete version of a file. A deletion process can set a deletion flag for a particular file. Heartbeat monitoring can detect that a number of files stored on a failing storage node need to be replicated to a new storage node.
Each storage node monitors the need for replication of the files it stores and maintains a replication list 55. The replication list 55 thus contains a number of files that need to be replicated. The files can be arranged in accordance with the priority of each replication. Typically, there can be three different priority levels. The highest level is reserved for files whose storage node has the last online copy. Such a file needs to be replicated quickly to other storage nodes so that a reasonable level of redundancy can be achieved. An intermediate level for the priority may refer to files where the versions do not match the storage nodes. A lower level of priority may refer to files stored on a failing storage node.
The storage node handles the files on the replication list 55 in accordance with their priority levels. The replication process is now described for a storage node
533 007, which is called the operating storage node here, although all storage nodes can function in this way.
The replication portion 53 of the maintenance process begins with the operating storage node trying to 71 become the "master" of the file it is trying to replicate. The operating storage node sends a unicast request to become the master of the other storage nodes known to store the file in question. The register 19 (cf. Fig. 1) provides a host list containing information regarding which storage nodes to request. In the event, for example, in the case of a collision request, that one of the storage nodes does not respond in the affirmative, the file is moved back to the list for the moment, and an attempt is made instead with the next file on the list. Otherwise, the operating storage node is considered "master" for this file and the other storage nodes set a flag indicating that the operating storage node is the master of the file in question.
The next step is to find 73 all copies of the file in question in the distributed storage system. This can be done by sending the operating storage node a multicast request to all storage nodes and asking which of them has the file. The storage nodes that have the file send replies to the request, which response contains both the version of the file they have and their host lists, ie the list of storage nodes that have the relevant file and that is stored in each storage node's register. These host lists are then merged 75 by the operating storage node to create a "master" host list, which corresponds to the Union of all the hosted host lists. If additional storage nodes are found that were not requested when the operating storage node made the attempt to become a master, this step can now be repeated for the additional storage nodes. Said master host list contains information regarding the versions of the file that different storage nodes have and gives an image of the file's status over the entire storage system.
If the operating storage node does not have the latest version of the file in question, 77 this file is obtained from one of the storage nodes that have the latest version.
The operating storage node then decides 79 whether the host list needs to be changed, typically if additional storage nodes should be added. If so, the operating storage node can perform a process as in a lot
533 007 is similar to the write process as performed by the server and as described in connection with Figures 4A-4C, and 5. The result of this process is that the file is written to a new storage node.
In the case of version conflicts, the operating storage node can update 81 copies of the file stored on other storage nodes so that all stored files have the correct version.
Excess copies of the stored file can be deleted 83. If the replication process is initiated by a deletion operation, the process can skip directly to this step. Then, as soon as all storage nodes have accepted the deletion of the file, simply the operating storage node, using unicast, requests that all storage nodes physically delete the file in question. The storage nodes confirm that the file has been deleted.
Furthermore, the status is updated, ie the file's master host list. Thereafter, it is optionally possible to repeat steps 73-83 to ensure that the need for replication no longer exists. This repetition will result in a uniform master host list that does not need to be updated in step 85.
Then, the replication process for this file is terminated and the operating storage node can release 87 the status of the file master by sending a corresponding message to all other storage nodes in the host list.
This system, where each storage node takes responsibility for the maintenance of all files it stores over the entire set of storage nodes, provides a self-repair (in the case of a failing storage node), self-cleaning (in the case of file conflicts or files to be removed) system with excellent reliability. It is easily scalable and can store files for a large number of applications simultaneously.
The invention is not limited to the particular examples shown and can be varied and modified in various ways within the scope of the appended claims.
533 007
9 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8 Sheet 9
40 members in 15 offices
Priority claims2
| Document | Office | Kind | Date |
|---|---|---|---|
| 0802277 | Sweden | A | |
| SE20080002277 | – | – | – |
Members40
| Document | Office | Kind | |
|---|---|---|---|
| SE0802277A1 | Sweden | A1 | |
| AU2009306386A1 | Australia | A1 | |
| CA2741477A1 | Canada | A1 | |
| WO2010046393A2 | World Intellectual Property Organization (WIPO) | A2 | |
| SE533007C2This record | Sweden | C2 | |
| EP2342663A2 | European Patent Office (EPO) | A2 | |
| KR20110086114A | Republic of Korea | A | |
| MX2011004240A | Mexico | A | |
| US2011295807A1 | United States of America | A1 | |
| CN102301367A | China | A | |
| ZA201103754B | South Africa | B | |
| US2012023179A1 | United States of America | A1 | |
| EA201100545A1 | Eurasian Patent Organization (EAPO) | A1 | |
| US2012030170A1 | United States of America | A1 | |
| JP2012506582A | Japan | A | |
| WO2010046393A3 | World Intellectual Property Organization (WIPO) | A3 | |
| US8688630B2 | United States of America | B2 | |
| US2014149351A1 | United States of America | A1 | |
| JP5553364B2 | Japan | B2 | |
| CN102301367B | China | B | |
| AU2009306386B2 | Australia | B2 | |
| EP2342663B1 | European Patent Office (EPO) | B1 | |
| US9026559B2 | United States of America | B2 | |
| DK2342663T3 | Denmark | T3 | |
| ES2538129T3 | Spain | T3 | |
| EP2908257A1 | European Patent Office (EPO) | A1 | |
| US9329955B2 | United States of America | B2 | |
| BRPI0914437A2 | Brazil | A2 | |
| KR101635238B1 | Republic of Korea | B1 | |
| US2016212212A1 | United States of America | A1 | |
| US9495432B2 | United States of America | B2 | |
| CA2741477C | Canada | C | |
| US2019266174A1 | United States of America | A1 | |
| EP2908257B1 | European Patent Office (EPO) | B1 | |
| EP3617897A1 | European Patent Office (EPO) | A1 | |
| US10650022B2 | United States of America | B2 | |
| EP3617897B1 | European Patent Office (EPO) | B1 | |
| US11468088B2 | United States of America | B2 | |
| US2023013449A1 | United States of America | A1 | |
| US11907256B2 | United States of America | B2 |
1 legal event, as the office reported them to INPADOC
Events
| Event | Code | |
|---|---|---|
| Patent has lapsedLapsedNUG | NUG |
Numbers
- Publication, DOCDB
- 533007
- Publication, EPODOC
- SE533007
- Application
- 802277
- Application, DOCDB
- 0802277
- Application, EPODOC
- SE20080002277
Titles2
- English
- Distributed data storage
- Swedish
- Distribuerad datalagring
Classification
- CPC, 8
- G06F11/1662
- G06F16/27
- G06F15/17331
- G06F11/2094
- H04L67/1097
- H04L69/40
- G06F16/1844
- H04L67/1095
- IPC, 2
- G06F15 173
- H04L69 40
