Distributed data storage.
Abstract
The present invention relates to a distributed data storage system comprising a plurality of storage nodes. Using unicast and multicast transmission, a server application may write data in the storage system. When writing data, at least two storage nodes are selected based in part on a randomized function, which ensures that data is sufficiently spread to provide efficient and reliable replication of data in case a storage node malfunctions.

Term
No projected expiry on record.
- Priority
- Filed
- Granted
- Today
9 claims: 7 independent, 2 dependent
- 1Claims Reivindicaciones 1. Un método para la escritura de datos en un sistema de almacenamiento de datos que comprende una pluralidad de nodos de memorización de datos, utilizándose dicho método en un servidor que ejecuta una aplicación que accede a datos en el sistema de almacenamiento de datos y que comprende:one. A method for writing data to a data storage system comprising a plurality of data storage nodes, said method being used on a server running an application that accesses data on the data storage system and comprising: el envío (41) de una consulta de memorización de multidifusión a una pluralidad de dichos nodos de memorización;sending (41) a multicast store query to a plurality of said store nodes;- la recepción (43) de una pluralidad de respuestas desde un subconjunto de dichos nodos de memorización, incluyendo dichas respuestas información del nodo de memorización que se relaciona, respectivamente, con cada nodo de memorización;- receiving (43) a plurality of responses from a subset of said storage nodes, said responses including information from the storage node that relates, respectively, to each storage node;- la selección (45) de al menos dos nodos de memorización en el subconjunto, en función de dicha respuesta, en donde la selección comprende: - the selection (45) of at least two storage nodes in the subset, depending on said response, where the selection comprises: - la determinación, basada en un algoritmo, para cada nodo de memorización en el subconjunto, de un factor de probabilidad que es función de su información del nodo de memorización y - the determination, based on an algorithm, for each storage node in the subset, of a probability factor that is a function of its information of the storage node and - la selección aleatoria de dichos al menos dos nodos de memorización, en donde la probabilidad de que se seleccione un nodo de memorización depende de su factor de probabilidad y - the random selection of said at least two storage nodes, where the probability that a storage node is selected depends on its probability factor and - el envío (47) de datos y de un identif icador de datos, correspondiente a los datos, a los nodos de memorización seleccionados. - sending (47) data and a data identifier, corresponding to the data, to the selected storage nodes. Geographical position includes the latitude and longitude of the memory node in question. posición geográfica incluye la latitud y la longitud del nodo de memorización en cuestión.
- 46. A method according to any of the preceding claims, wherein the information of the storage node includes the system load for the storage node in question. 6. Un método según cualquiera de las reivindicaciones precedentes, en donde la información del nodo de memorización incluye la carga del sistema para el nodo de memorización en cuestión.
- 57. A method according to any of the preceding claims, wherein the multicast memorization operational query includes a data identifier, which identifies the data to be memorized. 7. Un método según cualquiera de las reivindicaciones precedentes, en donde la consulta operativa de memorización de multidifusión incluye un identificador de datos, que identifica los datos que se van a memorizar.
- 68. A method according to any of the preceding claims, wherein at least three nodes are selected. 8. Un método según cualquiera de las reivindicaciones precedentes, en donde al menos se seleccionan tres nodos.
- 79. A method according to any of the preceding claims, wherein a list of storage nodes that successfully store the data is sent to the selected storage nodes. 9. Un método según cualquiera de las reivindicaciones precedentes, en donde se envía una lista de nodos de memorización que memorizan satisfactoriamente los datos a los nodos de memorización seleccionados.
- 810. A method according to any of the preceding claims, wherein said random selection is performed for a fraction of the nodes in the subset, whose fraction includes storage nodes with the highest probability factors. 10. Un método según cualquiera de las reivindicaciones precedentes, en donde dicha selección aleatoria se realiza para una fracción de los nodos en el subconjunto, cuya fracción incluye nodos de memorización con los más altos factores de probabilidad.
- 911. Un servidor adaptado para la escritura de datos en un sistema de memorización de datos que comprende una pluralidad de nodos de memorización de datos, comprendiendo dicho servidor:eleven. A server adapted for writing data to a data storage system comprising a plurality of data storage nodes, said server comprising: means for sending a multicast memorization operational query to a plurality of said memorization nodes;medios para enviar una consulta operativa de memorización de multidifusión a una pluralidad de dichos nodos de memorización;means for receiving a plurality of responses from a subset of said storage nodes, said responses including information from the storage node, respectively, in relation to each storage node;medios para la recepción de una pluralidad de respuestas desde un subconjunto de dichos nodos de memorización, incluyendo dichas respuestas información del nodo de memorización, respectivamente, en relación con cada nodo de memorización;means for selecting at least two memorization nodes in the subset, based on said responses, where the selection includes: medios para seleccionar al menos dos nodos de memorización en el subconjunto, en función de dichas respuestas, en donde la selección incluye: - la determinación, basada en un algoritmo, para cada nodo de memorización en el subconjunto, de un factor de probabilidad que es función de su información del nodo de memorización y - the determination, based on an algorithm, for each storage node in the subset, of a probability factor that is a function of its information of the storage node and - la selección aleatoria de dichos al menos dos nodos de memorización, en donde la probabilidad de que se seleccione un nodo de memorización depende de su factor de probabilidad y - the random selection of said at least two storage nodes, where the probability that a storage node is selected depends on its probability factor and - means for sending data and a data identifier, corresponding to the data, to the selected storage node. - medios para el envío de datos y de un identificador de datos, correspondiente a los datos, al nodo de memorización seleccionado.
Independent claims7
103 paragraphs in 1 section, as filed
(54) Title: STORAGE OF DISTRIBUTED DATA. (54) Title: DISTRIBUTED DATA STORAGE.
(57) Summary
The present invention relates to a distributed data storage system comprising a plurality of storage nodes. Using unicast and multicast transmission, a server application can write data to the memorization system. When writing data, at least two store nodes are selected based, in part, on a randomized function, which ensures that the data is sufficiently dispersed to provide efficient and reliable data replication in the event of abnormal operation of the memorization node.
(57) Abstract
The present invention relates to a distributed data storage system comprising a plurality of storage nodes. Using unicast and multicast transmission, a server application may write data in the storage system. When writing data, at least two storage nodes are selected based in part on a randomized function, which ensures that data is sufficiently spread to provide efficient and reliable replication of data in case a storage node malfunctions.
STORAGE OF DISTRIBUTED DATA
Technical field
The present invention relates to a method for writing data to a data storage system comprising a plurality of data storage nodes, the method of which is used on a server in the data storage system. The patent also refers to a server capable of carrying out the method.
Background
Such a method is disclosed, by way of example, in United States Patent 2005/0246393, Al. This method is disclosed for a system using a plurality of memory centers in disparate geographic areas. Distributed object storage managers are included to maintain information regarding memorized data.
One problem associated with such a system is how to perform simple, yet robust and reliable writing, and data maintenance.
Summary of the invention
An object of the present invention is therefore to securely write data to a distributed storage system.
The object of the invention is also achieved by means of a method for writing data to a data storage system of the initially mentioned class, which is performed on a server running an application that accesses data on the storage system. of data. The method comprises: sending a multicast memorization interrogation to a plurality of memorization nodes, receiving a plurality of responses from a subset of said memorization nodes, said responses including the information of memorization nodes, respectively, in relationship with each memorization node and the selection of at least two memorization nodes in the subset, depending on said responses. The selection includes the determination, based on an algorithm, for each storage node in the subset, of a probability factor that is a function of its storage node information and the selection, at random, of said at least two data nodes. memorization, where the probability that a memorization node is selected depends on its probability factor. The method further involves sending data and a data identifier, corresponding to the data, to the selected storage nodes.
This method performs a secure writing of data, since even when memorizing nodes are selected depending on their temporal aptitude, the information will continue to be dispersed, to a certain extent, through the system even for a short time interval. This means that the maintenance of the memory system will be less demanding, since the correlation of which memory nodes support the same information can be reduced to some extent. This means that a replication process, which can be performed when a storage node has an operational failure, can be carried out by a greater number of other storage nodes and consequently much faster. In addition, the risk of overloading the memory nodes with a high level during intensive write operations is reduced, since more memory nodes are used for writing and fewer are in idle mode.
The storage node information may include geographic data regarding the geographic position of each storage node, such as its latitude, longitude, and altitude. This allows the server to spread the information geographically, within a venue, a building, a country or on a global scale.
It is possible to allow the random selection of store nodes to be performed from store nodes in the subset that satisfy a primary criterion based on geographic separation, since this is an important feature for redundancy.
The information of the storage node can include the age of the system and / or the load of the system for the storage node in question.
The multicast memorization query may include a data identifier, which identifies the data to be memorized.
At least three nodes can be selected and a list of store nodes, which successfully store the data, can be sent to the selected store nodes.
The random selection of the memory nodes can be performed for a fraction of the nodes in the subset, which includes memory nodes with the highest probability factors. Thus, less suitable storage nodes are excluded, providing a selection of more reliable storage nodes while maintaining the random distribution of the information being written.
The present invention further relates to a server for writing data, corresponding to the method. The server, in such a case, usually includes means to carry out the actions of the method.
Brief description of the drawings
Figure 1 illustrates a distributed data storage system.
Figures 2A-2C and Figure 3 illustrate a process of reading data.
Figures 4A to 4C and Figure 5 illustrate a data writing process.
Figure 6 schematically illustrates a situation where multiple files are memorized between various data storage nodes.
Figure 7 illustrates the transmission of supervisory pulse signals.
Figure 8 is an overview of a data maintenance process.
Detailed description
The present invention 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 is represented in Figure 1.
A user computer 1 accesses, via the Internet 3, an application 5 running on a server 7. The user context, as illustrated herein, is therefore a regular client-server configuration, which is well known by itself. However, it should be noted that the data storage system disclosed by the invention may also be useful in other configurations.
In the illustrated case, two applications 5, 9 are running on server 7. Of course, however, this number of applications may be different. Each application has an API (Application Programming Interface) 11 that provides an interface in relation to the distributed data storage system 13 and supports demands, usually write and read demands, from applications running on the server. From an application point of view, the read or write information from / to the data storage system 13 need not appear different from the use of any other type of memorization solution, for example, a file server or just a hard drive.
Each API 11 communicates with the 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 will not be explained in more detail.
It should be noted that different API interfaces 11, on the same server 7 can access different sets of storage nodes 15. It should also be noted that there may be more than one server 7 that accesses each storage node 15. This operational circumstance, however, it does not further affect the way the memorization nodes operate, as will be described later.
The components of the distributed data storage system are the storage nodes 15 and the APIs 11 interfaces, in the server 7 that accesses the storage nodes 15. Therefore, the present invention refers to methods carried out in the server 7 and in the memorization nodes 15. Said methods will be mainly carried out as software implementation that run on the server and the storage nodes, respectively, and are determined, together, for the operation and properties of the global distributed data storage system.
The storage node 15 can usually be carried out by means of a file server that is provided with several functional blocks. The storage node can thus comprise a storage means 17, which usually includes several hard disk drives, optionally configured as a RAID system (Redundant Set of Independent Disks). Other types of memorizing means are, however, also possible.
The storing node 15 may further include a directory 19, comprising lists of data entity relationships / storing nodes as a hub list, as will be described later.
In addition to the hub list, each store node further contains a list of nodes that includes the IP addresses of all store nodes in their store or group of store nodes. The number of store nodes in a group can range from a few to hundreds of store nodes. The node list can also have a version number.
Also, the memorization node 15 may include a replication block 21 and an operational cluster monitor block 23. Replication block 21 includes a memorization node API 25 and is configured to execute functions to identify the need and perform a replication process, as will be described in detail later. The API interface of the memorization node 25 of the replication block 21 may contain a code that largely corresponds to the API interface code of the server memorization node 11
7, since the replication process comprises actions that correspond, to a large extent, to the actions performed by the server 7 during the read and write operations described below. As an example, the write operation performed during replication largely corresponds to the write operation performed by server 7. The operational grouping monitor block 23 is configured to perform monitoring of other storage nodes in the data storage system 13, as will be described in more detail below.
The storage nodes 15 of the distributed data storage system can be considered to exist at the same hierarchical level. There is no need to designate any master store node that is responsible for maintaining a directory of stored data entities and for monitoring data consistency, etc. Instead, all the storage nodes 15 can be considered the same and can sometimes perform the data management operations against other storage nodes in the system. This equality ensures that the system is operationally safe. In the event of an operational failure of the storage node, other nodes in the system will cover the malfunctioning node and ensure reliable data storage.
The operation of the system will be described in the following order: data reading, data writing and data maintenance. Although these methods work very well together, it should be noted that they can also, in principle, be done independently of each other. That is, by way of example, the data reading method can provide excellent properties even when the data writing method of the present invention is not used and vice versa.
The reading method is described below with reference to Figures 2A-2C and 3, the latter being a flow chart illustrating the method.
Reading, as well as other functions in the system, use multicast communication to communicate simultaneously with a plurality of memorization nodes. By means of an IP multicast or multicast, it is understood, in this case, a point-to-multipoint communication that is carried out by sending a message to an IP address that is reserved for multicast applications.
As an example, a message, usually a request, is sent to that IP address (eg, 244.0.0.1) and multiple receiving servers are registered as subscribers to that IP address. Each of the receiving servers has its 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 register as a subscriber for a multicast address, in which case point-to-point communication is achieved. However, within the context of this invention, such communication is nevertheless considered a multicast communication since a multicast system is used.
Unicast communication is also used with reference to communication with a single receiver.
Referring to Figure 2A and Figure 3, the method for retrieving data from a data storage system comprises sending 31 a multicast interrogation to a plurality of storage nodes 15. In the illustrated case, there are five Memorization nodes that each have an Internet Protocol (IP) address 192.168.1.1, 192.168.1.2, etc. The number of storage nodes is, of course, only by way of example. The interrogation contains, a data identifier 2B9B4A97-76E5-499E-A21A6D7932DD7927, which may be, by way of example, a Universally Unique Identifier, UUID, which is well known by itself.
The storage nodes themselves explore the data corresponding to the identifier. If such data is found, a store node sends a response, which is received 33 by server 7, see Figure 2B. As illustrated, the response may optionally contain additional information in addition to an indication that the store node has a copy of the relevant data. More specifically, the response may contain information from the store node directory about other store nodes containing the data, information regarding which version of the data is contained in the store node, and information regarding what load it is exposed to. the memorization node in the current situation.
Based on these responses, the server selects one or more store nodes from the data to be retrieved and sends 37 a request for data unicast to those store nodes, see Figure 2C.
In response to the data demand, the store node / nodes send the relevant data by unicasting to the server receiving the data. In the illustrated case, only one storage node is selected. Although this is sufficient, it is possible to select more than one storage node in order to receive two data sets that make a consistency check possible. If the data transfer fails, the server can select another store node to perform the recovery.
The selection of memorization nodes can be based on an algorithm that takes several factors into account in order to achieve adequate overall system performance. Under normal conditions, the memory node that has the latest data version and the lowest load will be selected even if other concepts are fully allowable.
Optionally, the operation can be concluded by the server by sending a list to all the storage nodes involved, indicating which nodes contain the data and with which version. On the basis of this information, the memorization nodes can themselves maintain the data properly through the replication process to be described.
Figures 4A-4C and Figure 5 illustrate a data writing process for the distributed data storage system.
Referring to Figure 4A and Figure 5, the method comprises a server that sends 41 a multicast store query to a plurality of store nodes. The memory query comprises a data identifier and basically consists of a query as to whether the receiving memory nodes can memorize this file. Optionally, the storage nodes can check their internal directories if they already have a file with this name and can notify server 7 in the unlikely event that this is the case, so that the server can rename the corresponding file.
In any case, at least a subset of the store nodes will provide response by unicast transmission to server 7. Under normal conditions, store nodes having a predetermined minimum free disk space will respond to this operational query. The server 7 receives 43 the responses comprising the information of the storage nodes in relation to the properties of each storage node, such as geographic data related to the geographical location of each server. By way of example, as indicated in Figure 4B, such geographic data may include the latitude, longitude, and altitude of each server. Other types of geographic data may, however, also be admissible, such as a ZIP postal code, a location string (ie, building, compound, rack row, rack column) or similar data.
As an alternative, or in addition to the geographic data, additional information related to the properties of the storage nodes can be provided, which serves as an input to a storage node selection process. In the illustrated example, the amount of free space at each storage node is provided along with an indication of the age of the storage node system and an indication of the load currently being experienced by the storage node.
Based on the responses received, the server selects at least two, in a typical embodiment three, store nodes in the subset to store the data. The selection of memorization nodes is carried out by means of an algorithm that takes into account different data. Selection can be made in order to achieve some kind of geographic diversity. At least it could be avoided, in a preferred embodiment, that only file servers in the same rack are selected as storage nodes. Under normal conditions, great geographic diversity can be achieved, even by selecting memorization nodes on different continents. In addition to geographic diversity, other parameters can be included in the selection algorithm. It is advantageous to have a randomized characteristic in the selection process as described below.
Under normal conditions, selection can be started by selecting several memorization nodes that are geographically sufficiently separated. This operation can be done in several ways. By way of example, there may be an algorithm that identifies multiple groups of store nodes, or store nodes may have group numbers, so that a single store node in each group can be easily captured.
The selection can then include calculating, on the basis of the information of memorization nodes of each node (system age, system load, etc.), a probability factor that corresponds to a proficiency score of the nodes of memorization. An older system, for example, that is less likely to have an operational failure, scores higher. The probability factor can thus be calculated as a dot product of two vectors, where one vector contains the information parameters of the storage nodes (or where applicable, their inverses) and the other contains the weighting parameters. corresponding.
The selection then comprises randomly selecting storage nodes, where the probability that a specific storage node is selected depends on its probability factor. Under normal conditions, if a first server has a probability factor twice as high as a second server, the first server has a double higher probability to be selected.
It is possible to remove a percentage of the store nodes with the lowest probability factors before random selection, so that this selection is made for a fraction of the nodes in the subset, whose fraction includes store nodes with the highest high probability factors. This is particularly useful if there are numerous memory nodes available that can take over the computation of the time consuming selection algorithm.
Of course, the selection process can be done in a different way. As an example, it is possible to first calculate the probability factor for all store nodes in the responding subset and perform random selection. When this operation is performed, it can be verified that the resulting geographic diversity is sufficient and, if it is not sufficient, repeat the selection with one of the two closest selected memorization nodes, excluded from the subset. Carrying out a first selection based on geographical diversity, for example, the acquisition of a memorization node in each group for subsequent selection based on the other parameters, is especially useful, again, in cases where There are numerous memorization nodes available. In these cases, a proper selection will continue to be made without the need to perform parameterized calculations of all available storage nodes.
When the memorization nodes have been selected, the data to be memorized and a corresponding data identifier are sent to each selected node, typically using a unicast transmission.
Optionally, the operation can be concluded by each storage node, which has successfully completed the write operation, by sending a confirmation to the server. The server then sends a list to all involved storage nodes indicating which nodes have successfully written data and which nodes have not done so properly. Based on this information, the storage nodes will be able to maintain the data on their own through the replication process described below. By way of example, if writing to a store node fails, there is a need for replication of the file to one or more store nodes in order to achieve the desired number of store nodes for that file.
The data writing method itself allows an API on a server 7 to memorize data in a very secure way, when excellent geographic diversity can be provided.
In addition to write and read operations, the API on server 7 can perform operations that delete files and update files. These processes will be described in relation to the data maintenance process later in this description.
The objective of the data maintenance process is to ensure that a reasonable number of storage nodes, without abnormal operation, each store the most recent version of each file. Furthermore, they can provide the function that no deleted file is memorized in any memorization node. Maintenance is carried out by the storage nodes themselves. Therefore, there is no need for a dedicated teacher to take responsibility for maintaining data memorization. This guarantees better reliability since the master would otherwise be a weak point in the system.
Figure 6 schematically illustrates a situation where multiple files are memorized between various data storage nodes. In the illustrated case, twelve nodes, which have consecutively numbered IP addresses from 192.168.1.1 to 192.168.1.12, are shown for illustration purposes. Of course, however,. IP address numbers need not be at the same level at all. The nodes are arranged in a circular order only to simplify the description, that is, the nodes need not have any particular order. Each node memorizes one or two files identified, for simplicity, by the letters AF.
Referring to Figure 8, the method for data maintenance comprises detecting 51 conditions in the data storage system that involve the need for data replication between nodes in the data storage system and a process of replication 53. The result of the detection process 51 is a list 55 of files for which a need for replication has been identified. The list may also include data regarding the priority of different replication needs. Based on this list, replication process 53 is performed.
The operational soundness of distributed storage lies in the fact that a reasonable number of copies of each file, in correct versions, are stored in the system. In the illustrated case, three copies of each file are memorized. However, if, by way of example, the storage node with the address 192.168.1.5 fails, the desired number of stored copies will not be reached for files B and
C.
An operational circumstance that gives rise to a need for replication is, therefore, the abnormal operation of a storage node in the system.
Each store node in the system can monitor the operational status of other store nodes in the system. This can be done by allowing each memorization node to emit a so-called supervision pulse signal, at periodic intervals, as illustrated in Figure 7. In the illustrated case, the storage node with the address 192.168.1.7 emits a multicast signal 57 to the other storage nodes in the system, indicating that it is working correctly. This signal may be received by all other operating memory nodes in the system that perform supervisory pulse signal supervision 59 (see Figure 8) or one of its subsets. In the case of the storage node with address 192.168.1.5, however, this node is operating abnormally and does not emit a supervisory pulse signal. Therefore, the other memorization nodes will note that no supervisory pulse signal has been emitted by this node for a long period of time, indicating that the memorization node in question is inactive.
The monitoring pulse signal may, in addition to the address of the storage node, include its version number from the node list. Another memorization node, which listens to the supervisory pulse signal and finds out that the transmitting memorization node has a list of subsequent memorization nodes, may then demand that the transmitting memorization node perform the transfer of its node list. This means that adding and removing store nodes can be accomplished by simply adding or deleting a store node and sending a new version of the node list to a single store node. This list of nodes will then be broadcast to all other store nodes in the system.
Again referring to Figure 8, each storage node searches 61 its internal directory for the files that are stored by the storage node in abnormal operation. The memory nodes that memorize files B and C by themselves will find the memory node in anomalous operation and therefore can add the corresponding file to their lists
55.
The detection process may, however, also reveal other conditions that imply the need for a file to be replicated. Normally, these conditions can be inconsistent, that is, that one or more storage nodes have an obsolete version of the file. A delete operation also involves a replication process since this process can perform the actual physical deletion of the file. The server deletion operation then only needs to make sure that the store nodes set a delete flag for the file in question. Each node can therefore supervise the read and write operations carried out in the data storage system.
The information provided by the server 7 at the conclusion of the read and write operations, respectively, may indicate that a single store node contains an outdated version of a file (in the case of a read operation) or that a node The memory protocol did not perform a write operation successfully. In both cases, there is a need to maintain data through replication, so that the global objects 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 that there is a need for replication, that is, the delete 67 and update 69 processes, of which a brief explanation.
The deletion process is initiated by server 7 (see Figure 1). Similar to the read process, the server sends a multicast query query to all store nodes in order to find out which store nodes have data with a specific data identifier. The storage nodes themselves perform a scan of the data with the relevant identifier and respond by unicast transmission if they have the data in question. The response may include a list, from the store node directory, of other store nodes containing the data. The server 7 then sends a unicast request to the storage nodes that are considered to store the file, in the sense that the file is deleted. Each storage node establishes an indicator in relation to the file and that 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 described below.
The update process has a search function, similar to that of the delete process, and a write function, similar to that performed in the write process. The server sends a query multicast query to all store nodes in order to find out which store nodes have data with a specific data identifier. The memory no.2s themselves scan the data with the relevant identifier and respond by unicast transmission if they have the data in question. The response may include a list, from the store node directory, of other store nodes containing the data. The server 7 then sends a unicast request, which communicates to the storage nodes the need to update the data. Of course, the lawsuit contains the updated data. The storing nodes, which update the data, send a confirmation to the server, which responds by sending a unicast transmission containing a list of the storing nodes that successfully updated the data and storing nodes that did not. Of
<td>again, this list is</td><td>you can use</td><td>by</td><td>the process</td><td>of</td>
<td>maintenance.</td><td></td><td></td><td></td><td></td>
<td>Again doing</td><td>reference to</td><td>the</td><td>Figure 8,</td><td>the</td>
<td>read operations</td><td>63, writing</td><td> 65,</td><td colspan="2">deletion 67 and</td>
<td>update 69 may</td><td colspan="2">indicate all of them</td><td>that exists</td><td>a</td>
need for replication. The same applies to the monitoring of pulse signals 59. The global detection process 51 thus generates data regarding which files need replication. As an example, a read or update operation may reveal that a specific store node contains an outdated version of a file. A deletion process can set a deletion flag for a specific file. Supervision by supervisory pulse signals may reveal that several files, memorized in a memorizing node, malfunctioning, need to be replicated for a new memorizing node.
Each of the archive nodes monitors the need for replication for all files and stores and maintains a replication list 55. Replication list 55 therefore contains several files that need replication. The files can be arranged in correspondence with the priority for each replication. Under normal conditions, there can be three different priority levels. The highest level is reserved for the files that the storage node maintains with the last online copy. Such a file needs rapid replication to other storage nodes so that a reasonable level of redundancy can be achieved. A medium priority level can refer to files where the versions are inconsistent between the storage nodes.
A lower priority level can refer to files that are stored in a storage node that has an abnormal operation.
The storage node is related to the files in replication list 55 according to their priority level. The replication process is described below for a store node which is referred to herein as the store store node, although all store nodes can operate in this manner.
The replication part 53 of the maintenance process begins with the attempt of the working storage node 71 to become the master node for the file attempting its replication. The working storage nodes send a unicast request to become the master node for other storage nodes that are known to store the file in question. Directory 19 (see Figure 1) provides a list of hubs containing information regarding which store nodes to request. Assuming, by way of example, in the event of a collision demand, that one of the storage nodes does not respond in the affirmative, the file will be moved to the present list again and an attempt will be made instead. with the next file in the list. If not, the working storage node is considered as the master of this file and the other storage nodes establish an indicator that the working storage node is the master node for the file in question.
The next step is to find all 73 copies of the file in question on the distributed storage system. This operation can be performed by the memorization node in operation that sends a multicast query to all the memorization nodes, asking which of them have the file. The storage nodes, which have the file, present responses to the query, which contains the version of the file that it maintains as well as its lists of concentrators, that is, the list of storage nodes that contain the relevant file is kept in the directory of each memorization node. These hub lists are then merged 75 by the operating memorization node, so that a master hub list is formed corresponding to the merging of all the retrieved hub lists. If additional store nodes are found, which were not interrogated when the operating store node attempted to become a master node, that step can now be repeated for the additional store nodes. The list of master node hubs contains information regarding which versions of the file the different storage nodes maintain and illustrate the operational status of the file within the entire storage system.
If the working storage node does not have the latest version of the file in question, this file is then retrieved 77 from one of the storage nodes that have the latest version.
The memorizing node in operation then decides 79 if the hub list needs to be changed, usually if additional memorizing nodes should be added. If so, the memory node in operation can perform a process very similar to the write process performed by the server and as described in relation to Figures 4A-4C and 5. The result of this process is that write the file to a new storage node.
In case of version inconsistencies, the working storage node can authorize 81 copies of the file that are stored in other storage nodes, so that all the stored files have the correct version.
Superfluous copies of the memorized file can be deleted 83. If the replication process is started by a delete operation, the process can go directly to this stage. Then, as soon as all the store nodes have accepted the deletion of the file, the running store node simply requests, using unicast, all store nodes to physically delete the file in question. The storage nodes confirm that the file is deleted.
In addition to the operational state, that is, the master file hub list is updated. It is then, optionally, possible to repeat steps 73-83 to ensure that there is no longer a need for replication. This repetition should result in a consistent master hub list that does not need to be updated in the stage.
85.
Later, the replication process for that file is concluded, and the working storage node can release 87 the operating state as master of the file by sending a message corresponding to all other storage nodes in the hub list.
This system, where each storage node assumes the responsibility of maintaining all the files it stores through the entire set of storage nodes, provides a self-repair system (in case of abnormal operation of the storage node) and of autodepuración (in case of inconsistencies of files or files to delete) with excellent reliability. It is easily scalable and can memorize files for a large number of different applications simultaneously.
The invention is not restricted to the specific examples disclosed and can be varied and modified in different ways within the scope of protection of the appended claims.
It is noted that in relation to this date, the best method known by the applicant to put the aforementioned invention into practice is the one that is clear from the present description of the invention.
10 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8 Sheet 9 Sheet 10
29 members in 13 offices
Priority claims2
| Document | Office | Kind | Date |
|---|---|---|---|
| 10160910 | European Patent Office (EPO) | A | |
| 2011056317 | European Patent Office (EPO) | W |
Members29
| Document | Office | Kind | |
|---|---|---|---|
| CA2797053A1 | Canada | A1 | |
| WO2011131717A1 | World Intellectual Property Organization (WIPO) | A1 | |
| EP2387200A1 | European Patent Office (EPO) | A1 | |
| US2012084383A1 | United States of America | A1 | |
| AU2011244345A1 | Australia | A1 | |
| IL222630A0 | Israel | A0 | |
| CN102939740A | China | A | |
| MX2012012198AThis record | Mexico | A | |
| EA201291023A1 | Eurasian Patent Organization (EAPO) | A1 | |
| JP2013525895A | Japan | A | |
| KR20130115983A | Republic of Korea | A | |
| ZA201208755B | South Africa | B | |
| EP2387200B1 | European Patent Office (EPO) | B1 | |
| EP2712149A2 | European Patent Office (EPO) | A2 | |
| EP2712149A3 | European Patent Office (EPO) | A3 | |
| US8850019B2 | United States of America | B2 | |
| JP5612195B2 | Japan | B2 | |
| US2014379845A1 | United States of America | A1 | |
| AU2011244345B2 | Australia | B2 | |
| CA2797053C | Canada | C | |
| CN102939740B | China | B | |
| BR112012027012A2 | Brazil | A2 | |
| US9503524B2 | United States of America | B2 | |
| US2017048321A1 | United States of America | A1 | |
| IL222630A | Israel | A | |
| EA026842B1 | Eurasian Patent Organization (EAPO) | B1 | |
| US9948716B2 | United States of America | B2 | |
| KR101905198B1 | Republic of Korea | B1 | |
| EP2712149B1 | European Patent Office (EPO) | B1 |
3 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Grant or registrationFG | FG | |
| Change of company name or juridical statusHC | HC | |
| Change of company name or juridical statusHC | HC |
Numbers
- Publication
- 2012012198
- Application
- 2012012198
Titles2
- English
- DISTRIBUTED DATA STORAGE.
- Spanish
- ALMACENAMIENTO DE DATOS DISTRIBUIDOS.
Classification
- CPC, 9
- H04L67/1097
- G06F15/17331
- H04L12/1845
- H04L67/1095
- G06F11/2094
- G06F16/1844
- G06F16/1824
- G06F16/29
- G06F3/067
- IPC, 1
- H04L29 08