Distributed storage system and method
12 claims: 6 independent, 6 dependent
- 1格納対象のデータを識別する識別子であるテーブル識別子に対応させて、複製を特定するレプリカ識別子と、前記レプリカ識別子に対応したデータ構造の種類を特定するデータ構造情報と、指定されたデータ構造に変換して格納されるまでのタイマ情報である契機情報と、を、前記データ構造の種類の数に対応させて備えたデータ構造管理情報と、 前記テーブル識別子に対応して、前記レプリカ識別子と、前記レプリカ識別子に対応した1つ又は複数のデータ配置先のデータノード情報とを備えたデータ配置特定情報と、 を記憶管理する構造情報保持部を有する構造情報管理装置と、 前記データ構造管理情報と前記データ配置特定情報とを参照して、更新処理のアクセス先のデータノードを特定する手段と、 それぞれがデータ格納部を備え、ネットワーク結合される複数のデータノード と、 を備え、 前記データノードは、 更新対象のデータを、一旦、書き込みデータ保持用の中間構造に格納して 応答を返すアクセス受付・処理部 と、 前記データ構造管理情報を参照し、指定された更新契機に応答して、前記中間構造に保持されるデータを、前記データ構造管理情報で指定されたデータ構造に変換する処理を行うデータ構造変換部と、 を備えている、ことを特徴とする分散ストレージシステム。
- 2前記データノードへのアクセス頻度の履歴を記憶するアクセス履歴記録部を備え、 前記データノードで非同期に行われる前記目的のデータ構造への変換の契機となる契機情報を、前記アクセス履歴記録部に記録されたアクセス情報に基づき、可変させる手段を備えている 、請求項1記載の分散ストレージシステム。
- 3予め定められたテーブル単位でデータ配置先のデータノード、配置先のデータノードにおける目的のデータ構造を制御する手段を備えた請求項1又は2記載の分散ストレージシステム。
- 4それぞれがデータ格納部を備え、ネットワーク結合される複数のデータノードを備え、 データの更新要求に対して前記データの複製先のデータノードでは、 更新対象のデータを、一旦、書き込みデータ保持用の中間構造に格納し、受け取った前記更新要求とは非同期で、それぞれ目的のデータ構造に変換して前記データ格納部に格納し、 前記データノードへのアクセス頻度の履歴を記憶するアクセス履歴記録部を備え、 前記データノードで非同期に行われる前記目的のデータ構造への変換の契機となる契機情報を、前記アクセス履歴記録部に記録されたアクセス情報に基づき、可変させる手段を備え、 格納対象のデータを識別する識別子であるテーブル識別子に対応させて、複製を特定するレプリカ識別子と、前記レプリカ識別子に対応したデータ構造の種類を特定するデータ構造情報と、指定されたデータ構造に変換して格納されるまでのタイマ情報である契機情報と、を、前記データ構造の種類の数に対応させて備えたデータ構造管理情報と、 前記テーブル識別子に対応して、前記レプリカ識別子と、前記レプリカ識別子に対応した1つ又は複数のデータ配置先のデータノード情報とを備えたデータ配置特定情報と、 を記憶管理する構造情報保持部を有する構造情報管理装置と、 前記データ構造管理情報と前記データ配置特定情報とを参照して、更新処理及び参照処理のアクセス先を特定するデータアクセス部を備えたクライアント機能実現部と、 それぞれが前記データ格納部を備え、前記構造情報管理装置と前記クライアント機能実現部とに接続される複数の前記データノードと、 を備え、 前記データノードは、 前記クライアント機能実現部からのアクセス要求に基づき、更新処理を行う場合に、中間構造にデータを保持して前記クライアント機能実現部に応答を返すアクセス受付・処理部と、 前記データ構造管理情報を参照し、指定された更新契機に応答して、前記中間構造に保持されるデータを、前記データ構造管理情報で指定されたデータ構造に変換する処理を行うデータ構造変換部と、 を備えたデータ管理・処理部を有する、ことを特徴とする、分散ストレージシステム。
- 5前記アクセス履歴記録部に記録されたアクセス情報、又は、前記アクセス情報を加工して得た別のアクセス情報を用いて、前記構造情報保持部の前記データ構造管理情報の更新契機情報を変更するか否か判定し、 前記データ構造管理情報の更新契機情報を変更する場合、前記構造情報管理装置に通知する変更判定部を備え、 前記構造情報管理装置は、前記変更判定部からの前記更新契機情報の変更の通知を受け、前記データ構造管理情報の更新契機情報を変更する構造情報変更部を備えた、請求項4記載の分散ストレージシステム。
- 6前記アクセス履歴記録部に記録されたアクセス情報が、前記データ格納部からの読み出しアクセスと、前記中間構造へのデータの書き込みアクセスの頻度情報を含む、請求項 2 又は5記載の分散ストレージシステム。
- 7前記データノードにおいて、 前記アクセス受付・処理部が、 アクセス受付部、アクセス処理部を備え、 前記データノードの前記データ格納部は、構造別データ格納部を備え、 前記アクセス受付部は、 前記クライアント機能実現部からの更新要求を受け付け、前記データ配置特定情報においてレプリカ識別子に対応して指定されているデータノードに対して更新要求を転送し、 さらに前記アクセス履歴記録部にアクセス要求を記録し、 前記データノードの前記アクセス処理部は、 受け取った更新要求の処理を行い、前記データ構造管理情報の情報を参照して更新処理を実行し、その際、前記データ構造管理情報の情報から、前記データノードに対する前記更新契機情報が零の場合、更新データを、前記データ構造管理情報に指定されるデータ構造に変換して、前記構造別データ格納部に格納し、 前記更新契機が零でない場合、前記中間構造に、一旦、更新データを書き込み、処理完了を応答し、 前記アクセス受付部は、 前記アクセス処理部からの完了通知、又は、 前記アクセス処理部からの完了通知及びレプリカ先の各データノードからの完了通知、 を受けると、前記クライアント機能実現部に対して応答し、 前記データ構造変換部は、前記中間構造のデータを、前記データ構造管理情報に指定されているデータ構造に変換し変換先の前記構造別データ格納部に格納する、請求項5記載の分散ストレージシステム。
- 8前記目的のデータ構造が同一の少なくとも二つのデータノードを備え、 前記二つのデータノードでは、前記書き込みデータ保持用の中間構造に保持されたデータから前記目的のデータ構造への変換を、設定された前記契機情報に基づき、それぞれ、時間的に重ならないタイミングで行い、一方のデータノードで、前記書き込みデータ保持用の中間構造に保持されたデータを、前記目的のデータ構造に変換しているとき、他方のデータノードでは、前記目的のデータ構造に変換されたデータの読み出しが行われる、請求項 2又は4 記載の分散ストレージシステム。
- 9それぞれがデータ格納部を備え、ネットワーク結合される複数のデータノードを備えた分散ストレージのデータ複製方法において、 データの更新要求に対応したデータの複製にあたり、複製先のデータノードでは、 更新対象のデータを、一旦、書き込みデータ保持用の中間構造に格納し、更新要求とは非同期で、それぞれ目的のデータ構造に変換して前記データ格納部に格納し、 前記データノードで非同期に行われる前記目的のデータ構造への変換の実行の契機となる契機情報を、前記データノードへのアクセスの履歴情報に基づき、可変させ 、 格納対象のデータを識別する識別子であるテーブル識別子に対応させて、複製を特定するレプリカ識別子と、前記レプリカ識別子に対応したデータ構造の種類を特定するデータ構造情報と、指定されたデータ構造に変換して格納されるまでの時間情報である契機情報と、を、前記データ構造の種類の数に対応させて管理するデータ構造管理情報と、 前記テーブル識別子に対応して、前記レプリカ識別子と、前記レプリカ識別子に対応した1つ又は複数のデータ配置先のデータノード情報とを備えたデータ配置特定情報と、 を構造情報管理装置の構造情報保持部にて記憶し、 データアクセス部において、前記データ構造管理情報と前記データ配置特定情報とを参照して、更新処理及び参照処理のアクセス先を特定し、 前記データノードは、 クライアントからのアクセス要求に基づき、更新処理を行う場合に、中間構造にデータを保持して応答を返し、 前記データ構造管理情報を参照し、指定された更新契機に応答して、前記中間構造に保持されるデータを、前記データ構造管理情報で指定されたデータ構造に変換する、 ことを特徴とする、データ複製方法。
- 10予め定められたテーブル単位でデータ配置先のデータノード、配置先のデータノードにおける目的のデータ構造を制御する請求項 9 記載のデータ複製方法。
- 11前記アクセスの履歴情報が、前記データ格納部からの読み出しアクセスと、前記中間構造へのデータの書き込みアクセスの頻度情報を含む、請求項9記載のデータ複製方法。
- 12前記目的のデータ構造が同一の少なくとも二つのデータノードを用意し、 前記二つのデータノードでは、前記書き込みデータ保持用の中間構造に保持されたデータから前記目的のデータ構造への変換を、設定された前記契機情報に基づき、それぞれ、時間的に重ならないタイミングで行い、一方のデータノードで、前記書き込みデータ保持用の中間構造に保持されたデータを、前記目的のデータ構造に変換しているとき、他方のデータノードでは、前記目的のデータ構造に変換されたデータの読み出しが行われる、請求項9記載のデータ複製方法。
Independent claims12
142 paragraphs, as filed
0001(Description of related application) The present invention is based on the priority claim of Japanese Patent Application No. 2011-169588 (filed on August 2, 2011), and all the contents of the application are incorporated in this document by citation. It shall be. The present invention relates to distributed storage, and more particularly to distributed storage systems, methods and devices capable of controlling data structures.
0002Distribution that realizes a system that connects multiple computers (data nodes, or simply "nodes") to a network and stores and uses data in the data storage unit (HDD (Hard Disk Drive), memory, etc.) of each computer. A storage system (Distributed Storage System) is being used.
0003With common distributed storage technology, -Which computer (node) to place the data on -Which computer (node) should be used for processing? Such a judgment is realized by software or special dedicated hardware. In the distributed storage system, the resource usage in the system is adjusted by dynamically changing the operation according to the state of the system, and the performance for the system user (client computer) is improved.
0004In a distributed storage system, data is distributed to a plurality of nodes, so a client trying to access data first needs to know which node holds the data. In addition, a client trying to access data needs to know which node (one or more) to access when there are a plurality of nodes having the data.
0005In a distributed storage system, a method of separately storing a file body and metadata (file storage location, file size, owner, etc.) of the file is generally used for file management.
0006In a distributed storage system, the metaserver method is known as one of the technologies for a client to know the node holding data. In the meta server method, a meta server composed of one or more (however, a small number) of computers that manages the location information of data is provided. However, in the distributed storage system of the meta server method, the processing performance of the meta server that performs the process of detecting the position of the node holding the data is insufficient due to the increase in the scale of the system configuration (managed by one meta server). The number of nodes to be operated becomes enormous, and the processing performance of the metaserver cannot keep up), and the introduced metaserver may become a bottleneck in access performance.
0007<Distributed KVS> As another method (technique) for knowing the position of the node holding the data, there is a method of finding the position of the data by using a distribution function (for example, a hash function). This kind of method is, for example, distributed KVS (Key Value). Store: Key-value store). Distributed KVS is a type of distributed storage system that realizes a simple data model storage function consisting of pairs of "Key" and "Value" like an associative array on multiple nodes. is there. In a distributed storage system based on the distributed KVS method (also called a distributed KVS system), all clients share a distributed function and a list of nodes participating in the system (node list). In addition, the stored data is divided into fixed-length or arbitrary-length data fragments (Value). Each data fragment is given an identifier that can uniquely identify the data fragment, and the location of the data fragment is determined by using the identifier and the distribution function. For example, since the storage destination node (server) differs depending on the key value depending on the hash function, it is possible to store data in a distributed manner on a plurality of nodes. Further, if the distribution functions are the same, the storage destination based on the same key is always the same, so that the accessing client can easily grasp the data access destination. In a simple distributed KVS system, a data access function based on Key and Value is realized by using Key as an identifier and Value corresponding to Key as a unit of stored data.
0008In a distributed storage system based on the distributed KVS method, when each client accesses data, the key is used as the input value of the distributed function, and the position of the node storing the data is based on the output value of the distributed function and the node list. Is calculated mathematically.
0009In a distributed storage system based on the distributed KVS method, among the information shared between clients, the distributed function basically does not change over time (time-invariant). On the other hand, the contents of the node list are changed at any time due to the failure or addition of the node. Therefore, it is necessary for the client to be able to access the information in any way.
0010<Replication> In a distributed storage system, in order to ensure availability (ability of the system to operate continuously), it is common to hold data replication on multiple nodes and utilize data replication for load balancing. It has been.
0011Patent Document 1 discloses a technique for realizing load balancing by duplicating the created data. Further, in Patent Document 2, the server defines the information structure definition structure in the information structure definition part, the registration client constructs the database by the information structure definition structure, generates the database access tool, and uses this tool to generate the database. The configuration for registering information is disclosed in. Patent Document 3 also describes a storage node in a distributed storage system that stores a copy of an object that each copy can access via its own unique locator value, and a keymap that stores a keymap entry for each object. Each keymap entry, including an instance, for a given object discloses a duplicate of the object, the corresponding key value, and a configuration that includes each locator. Further, in Patent Document 4 (including the inventor of the present application as a co-inventor), every time the data is updated, the changed contents are saved in chronological order, and the data writing to the storage is tracked and captured, and the data is updated. When an occurrence occurs, the changes can be journalized to the secondary storage (change history database) to reproduce the data at any time in the past (Any Point In Time (APIT) Recovery), and the data. CDP (Continuous Data Protection; Continuous Data Protection) is disclosed. Patent Document 4 is a storage system equipped with a data protection function that makes it possible to restore data at a past point in time by recording the changed contents as a log in chronological order when data is updated. Then, based on the analysis result of the history information of access to the storage and / or the information notified from the outside, a predetermined opportunity regarding data access is extracted, and the data corresponding to the extracted predetermined opportunity is obtained. It is created from the data stored in the storage and the log information, and the created data is stored in the storage as data corresponding to the predetermined opportunity.
<p num="0012"><patcit num="1"><text>Japanese Unexamined Patent Publication No. 2006-12005 (Patent No. 4528039)</text></patcit><patcit num="2"><text>Japanese Unexamined Patent Publication No. 11-195044 (Patent No. 3911810)</text></patcit><patcit num="3"><text>Special Table 2009-522659 Gazette</text></patcit><patcit num="4"><text>Japanese Unexamined Patent Publication No. 2007-317017</text></patcit></p>
<p num="0013"> Each disclosure of each of the above patent documents is incorporated herein by reference. The following is an analysis of related technologies.</p><p num="0014"> In a distributed storage system of related technology, data replication is held by multiple nodes in order to maintain availability, but the same physical structure is held by multiple nodes. As a result, access response performance and availability are guaranteed in the distributed storage system. However, since the duplicated data is held in the same physical structure in a plurality of nodes, for example, in an application that reads and analyzes the data, the data is different from the data structure of the held duplicated data. For applications and the like used in a data structure, it is necessary to prepare a storage for converting to another data structure and holding another data structure. Conversion to another data structure causes an increase in processing load and processing delay, resulting in an increase in storage capacity for holding another data structure.</p><p num="0015"> At that time, the inventors of the present application hope that, for example, a special improvement in performance can be expected by taking special measures regarding the execution of writing (writing, updating) of the data and the conversion of the data into the target data structure. Since I found out, I will propose this this time.</p><p num="0016"> An object of the present invention is to provide a distributed storage system and a method capable of ensuring availability in data replication in distributed storage and improving both write performance and processing performance on the read side.</p>
<p num="0017"> According to the present invention, in order to solve at least one of the above problems, the configuration is roughly as follows (however, it is not limited to the following).</p><p num="0018"> According to the present invention, each of the data nodes is provided with a data storage unit, and a plurality of data nodes connected to the network are provided. In response to a data update request, the data node to which the data is duplicated temporarily writes the data to be updated and holds the data. It is stored in the intermediate structure for, and asynchronously with the update request, it is converted to the desired data structure and stored in the data storage unit. It is provided with an access history recording unit that stores the history information of access to the data node. Distributed, provided with means for varying the trigger information that triggers the execution of conversion to the target data structure that is performed asynchronously at the data node, based on the access history information recorded in the access history recording unit. A storage system is provided.</p><p num="0019"> According to the present invention, in data replication of distributed storage, each of which has a data storage unit and a plurality of network-coupled data nodes. When duplicating data in response to a data update request, the data node at the replication destination The data to be updated is temporarily stored in the intermediate structure for holding the write data, and asynchronously with the update request, each is converted into the target data structure and stored in the data storage unit. Provided is a data duplication method for distributed storage, which changes the trigger information that triggers the execution of conversion to the target data structure asynchronously performed at the data node based on the access history information of the data node.</p>
<p num="0020"> According to the present invention, it is possible to ensure availability in data replication in distributed storage and to improve both write performance and read-side processing performance.</p>
0021<figref num="1">It is a figure which shows the system structure of one Embodiment of this invention.</figref><figref num="2">It is a figure which shows the structural example of the data node of one Embodiment of this invention.</figref><figref num="3">It is a figure which shows typically the data structure management information 921 in one example embodiment of this invention.</figref><figref num="4">It is a figure which shows typically an example of the data holding structure of the table in one exemplary Embodiment of this invention.</figref><figref num="5">It is a figure which shows the example of the data arrangement specific information 922 in one example embodiment of this invention.</figref><figref num="6">It is a figure explaining data retention and asynchronous update schematically.</figref><figref num="7">It is a figure explaining the Write process and the analysis system process in FIG. 6 schematically.</figref><figref num="8">It is a figure which shows typically the data holding and asynchronous update in one example embodiment of this invention.</figref><figref num="9">It is a figure which shows the structural example of the access history recording part and the structural information management means of one Embodiment of this invention.</figref><figref num="10">It is a flowchart explaining the operation of the access processing in the client function realization means 61 in one exemplary Embodiment of this invention.</figref><figref num="11">It is a flowchart explaining the operation of the access processing in the data node in one exemplary embodiment of the present invention.</figref><figref num="12">It is a flowchart explaining the data conversion process in one exemplary Embodiment of this invention.</figref><figref num="13">It is a figure (the 1) explaining the operation sequence of the Write process in one example embodiment of this invention.</figref><figref num="14">It is a figure (the 2) explaining the operation sequence of the Write process in one example embodiment of this invention.</figref><figref num="15">It is a figure explaining another exemplary embodiment of this invention.</figref><figref num="16">It is a figure explaining still another exemplary embodiment of this invention.</figref>
0022Some preferred embodiments for carrying out the invention will be described. In some preferred embodiments, each has a data storage unit and a plurality of data nodes connected to the network. For example, when duplicating data at the time of data update, the replication destination data node once transfers the data to be updated. , Stored in an intermediate structure for holding write data (Queue (queue), FIFO (First In First Out), Log (log), etc.), asynchronous with the update request, converted to the desired data structure, and described above. Store in the data storage unit (12). Further, the data node includes an access history recording unit (71) that stores a history of access frequency to the data node. In the data node, the access history information (access frequency) stored in the access history recording unit (71) is the trigger information that triggers the execution of the conversion to the target data structure that is performed asynchronously in the data node. Set variably based on.
0023In some preferred embodiments, each of the replication destination data nodes holds the data in the intermediate structure, returns a response, and transfers the data structure held in the intermediate structure to the data to be updated. When the time specified in the trigger information elapses from the reception, the data may be asynchronously converted to the target data structure and then stored in the data storage unit.
0024In some preferred embodiments, the target data structure in the data node of the data allocation destination and the data node of the allocation destination may be controlled in a predetermined table unit.
0025In some preferred embodiments, a replica identifier that identifies a replica, a data structure information that identifies the type of data structure that corresponds to the replica identifier, and data structure information that corresponds to a table identifier that identifies the data to be stored. Data structure management information provided with trigger information, which is timer information until converted to a specified data structure and stored, corresponding to the number of types of the data structure (921 in FIG. 2: FIG. 3). When, Data allocation specific information (922 in FIG. 2: FIG. 5) including the replica identifier and data node information of one or a plurality of data allocation destinations corresponding to the replica identifier, corresponding to the table identifier. The access destination of the update process and the reference process is specified by referring to the structural information management device (9) having the structural information holding unit (92) that stores and manages the data, the data structure management information, and the data arrangement specific information. A plurality of client function realization units (61) having a data access unit, and a plurality of units each having the data storage unit (12) and connected to the structural information management device (9) and the client function realization unit (61). The data nodes (1 to 4) of the above are provided. When the data node performs update processing based on an access request from the client function realization unit (61), the data node temporarily holds data in the intermediate structure and then returns a response to the client function realization unit (61). The data held in the intermediate structure is the data specified in the data structure management information in response to the designated update opportunity by referring to the reception / processing unit (111, 112) and the data structure management information. It may be configured as a data management / processing unit (11) including a data structure conversion unit (113) that performs processing for converting into a structure.
0026In some preferred embodiments, the data structure of the structure information holding unit is used by using the access information recorded in the access history recording unit (71) or another access information obtained by processing the access information. When determining whether or not to change the update trigger information of the management information (921) and changing the update trigger information of the data structure management information (921), the change determination unit (72) notifying the structure information management device is used. The structural information management device (9) receives a notification of a change in the update trigger information from the change determination unit (72), and the structural information change unit (91) changes the update trigger information of the data structure management information. To be equipped. In a preferred embodiment, the access history recording unit (71) may record the access frequency as access information.
0027In some preferred embodiments, the access information recorded in the access history recording section (71) includes frequency information of read access from the data storage section and write access to the intermediate structure (or data write access). It may be information that indicates the pattern of access occurrence, the tendency of access occurrence, etc.).
0028In some preferred embodiments, the data node comprises an access reception unit (111), an access processing unit (112), and a data structure conversion unit (113). The data storage unit (12) of the data node includes structural data storage units (121 to 123), and the access reception unit (111) receives an update request from the client function realization unit and arranges the data. The update request is transferred to the data node specified corresponding to the replica identifier in the specific information, the access request is further logged in the access history recording unit, and the access processing unit (112) of the data node receives the data node. The update request is processed, and the update process is executed with reference to the information of the data structure management information. At that time, when the update trigger information for the data node is zero from the information of the data structure management information, the update data is converted into the data structure specified in the data structure management information, and the data by structure is stored. If the update trigger is not zero, the update data is once written in the intermediate structure, and the processing completion is replied. The access reception unit (111) Completion notification from the access processing unit (Fig. 14) or Completion notification from the access processing unit and completion notification from each data node of the replica destination (Fig. 13), When it receives, it responds to the client function realization unit (9) and responds. The data structure conversion unit (113) converts the data held in the intermediate structure into the data structure specified in the data structure management information, and converts the data into the structure-specific data storage units (121 to 123) of the conversion destination. It may be stored.
0029Some exemplary embodiments will be described below.
0030<System configuration> FIG. 1 is a diagram showing an example of a system configuration according to an exemplary embodiment of the present invention. It is equipped with data nodes 1 to 4, network 5, client nodes 6, and structural information management means (structural information management device) 9.
0031Data nodes 1 to 4 are data storage nodes that constitute distributed storage, and are composed of one or more arbitrary numbers. Network 5 realizes communication between network nodes including data nodes 1 to 4. Client node 6 is a computer node that accesses distributed storage. Client node 6 does not necessarily have to exist independently. An example in which data nodes 1 to 4 also serve as a client computer will be described later with reference to FIG. Data nodes 1 to 4 are data management / processing means (data management / processing unit) 11, 21, 31, 41, data storage units 12, 22, 32, 42, and access history recording units 71-1 to 71-, respectively. Equipped with 4.
0032The data management / processing means X1 (X = 1, 2, 3, 4) receives an access request to the distributed storage and executes processing. The data storage unit X2 (X = 1, 2, 3, 4) holds and records the data in charge of the data node.
0033The client node 6 includes a client function realization means (client function realization unit) 61. The client function implementation means 61 accesses the distributed storage composed of the data nodes 1 to 4. The client function realizing means 61 includes a data access means (data access unit) 611.
0034The data access means (data access unit) 611 acquires structural information (data structure management information and data arrangement specific information) from the structural information management means 9, and uses the structural information to specify the access destination data node.
0035In each data node 1 to 4 or any device (switch, intermediate node) in the network 5, part or all of the structural information stored in the structural information holding unit 92 of the structural information management means 9 is stored in the own device. Alternatively, it may be held in a cache (not shown) in another device.
0036The structural information stored in the structural information holding unit 92 may be accessed to a cache (not shown) arranged in the own device or in a predetermined predetermined location. Since a known distributed system technique can be applied to the synchronization of the structural information stored in the cache (not shown), the details will be omitted here. As is well known, storage performance can be accelerated by using cache.
0037The structural information management means (structural information management device) 9 includes a structural information changing means 91 for changing the structural information and a structural information holding unit 92 for holding the structural information. The structure information holding unit 92 includes data structure management information 921 (see FIG. 2) and data arrangement specific information 922 (see FIG. 4). The data structure management information 921, which will be described later with reference to FIG. 3, includes a replica identifier that specifies replication and a data structure information that specifies the type of data structure corresponding to the replica identifier with respect to the table identifier. , Has an entry consisting of an update trigger, which is time information until it is stored as a specified data structure, for the number of duplicates of data. The data allocation specific information 922, which will be described later with reference to FIG. 5, provides the replica identifier and the data node information of one or more data allocation destinations corresponding to the replica identifier corresponding to the table identifier. Have.
0038The access history recording units 71-1 to 4 record log information of Read access and Write access of data nodes 1 to 4. As the access log information, frequency information corresponding to the number of accesses within a predetermined period may be stored.
0039In FIG. 1, the client node 6 is provided independently (separately) from the data nodes 1 to 4, but the client node 6 is not necessarily provided independently (separately) from the data nodes 1 to 4. Not needed. That is, as will be described below as a modified example, any one or more of the data nodes 1 to 4 may be provided with the client function realizing means 61.
0040<Data node configuration example> FIG. 2 is a diagram for explaining the configuration example of FIG. 1 in detail. FIG. 2 shows the configuration centered on the data nodes 1 to 4 in FIG. Since the data nodes 1 to 4 in FIG. 1 have basically the same configuration, in FIG. 2, the data management / processing means 11, the data storage unit 12, and the access history recording unit 71 of the data node 1 (71- in FIG. 1). (Corresponding to 1) is shown. In drawings such as FIG. 2, for simplification, the structural information stored in the structural information holding unit 92 may be referred to by reference numeral 92.
0041The data management / processing means 11 of the data node 1 includes an access receiving means (access receiving unit) 111, an access processing means (access processing unit) 112, and a data structure conversion means (data structure conversion unit) 113. The data management / processing means 21, 31, and 41 of the other data nodes 2 to 4 have the same configuration.
0042The access receiving means 111 receives an access request from the data access means 611, and returns a response to the data access means 611 after the processing is completed.
0043The access processing means 112 uses the structural information of the structural information holding unit 92 (or cache information held at an arbitrary location thereof) to transfer the access processing to the corresponding data storage unit 12X (X = 1, 2, 3). Do it against.
0044The access receiving means 111 records the information of the access request (access command) in the access history recording unit 71 together with the reception time information, for example.
0045The data structure conversion means 113 converts the data of the structure-specific data storage unit 121 into the structure-specific data storage unit 12X (X = 1, 2, 3) at regular intervals.
0046The data storage unit 12 includes a plurality of types of structural data storage units. Although not particularly limited, FIG. 2 includes a structural data storage unit 121 (data structure A), a structural data storage unit 122 (data structure B), and a structural data storage unit 123 (data structure C). The type of data structure to be selected is arbitrary for each structure-specific data storage unit 12X (X = 1, 2, 3).
0047The structure-specific data storage unit 121 (for example, data structure A) has a structure specialized in response performance to processing (addition or update of data) involving writing of data. Specifically, software that holds data changes in queues (for example, FIFO (First In First Out)) on high-speed memory (dual port RAM (Random Access Memory), etc.), and arbitrary storage media for access request processing contents. Software etc. to be added as a log is implemented in. The data structure B and the data structure C are different data structures from the data structure A, and have different data access characteristics from each other. The data storage unit 12 does not necessarily have to be a single storage medium. The data storage unit 12 of FIG. 4 may be realized as a distributed storage system composed of a plurality of data arrangement nodes, and the data storage unit 12X for each structure may be distributed and stored.
0048The data arrangement specific information 922 is information (and means for storing and acquiring the information) for specifying the storage destination of the data to be stored in the distributed storage or the data fragment. As the data distribution method, for example, the meta server method or the distributed KVS method is used as described above.
0049In the case of the meta server method, the information for managing the position information of the data (for example, the block address and its corresponding data node address) is the data arrangement specific information 922. By referring to this information (metadata), the meta server can know where to place the necessary data.
0050In the case of the distributed KVS method described above, the list of nodes participating in the system corresponds to this data arrangement specific information. The data node of the data storage destination can be determined by using the identifier for storing the data and the node list information.
0051The data access means 611 should access the data nodes 1 to 4 using the data arrangement specific information 922 in the structural information management means 9 or the cache information of the data arrangement specific information 922 stored in a predetermined predetermined location. Is specified, and an access request is issued to the access receiving means 111 of the data node.
0052<Data structure management information> The data structure management information 921 in FIG. 2 is parameter information for specifying a data storage method for each set of data. FIG. 3 is a diagram showing an example of the data structure management information 921 of FIG. Although not particularly limited, in the example shown in FIG. 3, the unit for controlling the data storage method is a table. Then, for each table (for each table identifier), each information of the replica identifier, the type of the data structure, and the update trigger is prepared for the number of duplicates of the data duplication.
0053In Figure 3 (A), each table holds three copies (but the number of copies is not limited to three) to ensure availability (hold). The replica identifier is information that identifies each replica, and is assigned as 0, 1, and 2 in FIG. 3 (A). The data structure is information indicating a data storage method. In FIG. 3 (A), three types of data structures (A, B, C) are specified in different methods for each replica identifier.
0054Figure 3 (B) shows an example of data storage methods for data structures A, B, and C (but not limited to these storage methods). In the example of Fig. 3 (B), the type of data storage method is A: Queue, B: Low store, C: Column store Is specified. In the example of FIG. 3B, the replica identifier 0 of the table identifier "Stocks" is stored as the data structure B (low store).
0055Each data structure is a method for storing data. A: A queue is a Linked List.
0056B: ROW STORE stores the records in the table in ROW order.
0057C: Column store (COLUMN STORE) stores in column (COLUMN) order.
0058<Table configuration example> FIG. 4 is a diagram schematically showing an example of the data holding structure of the table. The table (A) in FIG. 4 has a Key column and three Value columns, and each row consists of a Key and a set of three Values.
0059The column store and the row store are formats in which the storage order on the storage medium is stored in a row (row) base and a column (column) base, respectively. As a storage method for the table (see (A) in Fig. 4), It is retained in the data structure B (low store) as the data of replica identifiers 0 and 1 (see (B) and (C) in Fig. 4). It is retained as data structure C (column store) as the data of replica identifier 2 (see (D) in Fig. 4).
0060<Update opportunity information> With reference to FIG. 3 (A) again, the update trigger in the data structure management information 921 (see FIG. 2) is the time trigger until the data is stored as the specified data structure. In the example of Stocks replica identifier 0, it is specified as 30 sec. Therefore, in the data node that stores the data structure B (low store) of the replica identifier 0 of Stocks, it is shown that the update of the data is reflected to the structure-specific data storage unit 122 of the low store method in 30 seconds. .. Until the data update is reflected, the data is retained as an intermediate structure such as a queue. In addition, the data node also responds to the request from the client by storing it in the intermediate structure. In this embodiment, the conversion to the specified data structure is performed asynchronously with respect to the update request.
0061In the following, the data to be updated is transferred between the data nodes in a synchronous manner, and the data structure is converted to the target structure asynchronously. An example in which a timer is used as update trigger information for asynchronously converting a data structure will be described (however, the present invention is not limited to the following implementation).
0062<Data placement specific information> FIG. 5 is a diagram showing an example of the data arrangement specific information 922 of FIG. A placement node (data node for data storage) is specified for each of the replica identifiers 0, 1, and 2 (see Fig. 3) of each table identifier. This corresponds to the metaserver method described above. In the case of the distributed KVS method, the data arrangement specific information 922 corresponds to the node list information (not shown) participating in the distributed storage. By sharing this node list information between data nodes, for example, the placement node can be specified by the consistent hashing method using "table identifier" + "replica identifier" as key information. In addition, the replica can be stored in the adjacent node in the consistent hashing method as the placement destination.
0063<Write intermediate structure: comparison example> FIG. 6 is a diagram schematically explaining the basic format of table data retention and asynchronous update. FIG. 6 is also a diagram for explaining a problem to be solved by the present invention, and therefore also a diagram for explaining a comparative example of the present invention.
0064When the value of the update trigger information is larger than 0, each data node has an intermediate structure (also called "Write priority structure" or "Write intermediate structure") that has an excellent response speed for Write (update request). And accept updates. Write When writing to the intermediate structure, a processing completion response is returned to the client that requested the update.
0065The update data written in the Write intermediate structure of each data node is updated asynchronously to the conversion target data structure in each data node. In the example shown in FIG. 6, in the data node whose replica identifier is 0 by Write, the data structure A is stored and held in the Write intermediate structure, and the data nodes with replica identifiers 1 and 2 are stored in a synchronous manner (Synchronous). , Write The data of the data structure A held in the intermediate structure is replicated. In each of the data nodes of replica identifiers 1 and 2, the data of data structure A transferred from the data nodes of replica identifiers 0 and 1, respectively, is temporarily stored and held in the Write intermediate structure. In the data nodes corresponding to the data structures corresponding to the replica identifiers 0, 1, and 2, the conversion to the target data structures B and C is the update trigger information of the data structure management information 921 as shown in FIG. 3 (A). Specified by. For example, in the data node with replica identifier 0, the timer is started from Write in the data structure A, and when 30 seconds (seconds) elapse (timeout: update trigger occurs), the data structure A is converted to the data structure B (Row-Store). In the data node with replica identifier 1, when the data structure A transferred by the synchronization method (Sync) is received from the data node with replica identifier 0, the timer is started, and when 60 seconds have passed (timeout: update trigger occurs). , Convert to data structure B (Row-Store). In the data node of replica identifier 2, when the data structure A transferred by the synchronization method (Sync) is received from the data node of replica identifier 1, the timer is started, and when 60 seconds have passed (timeout: update trigger occurs). , Convert to data structure C (Column-Store).
0066As shown in FIG. 6, replication between data nodes of update data (data structure A) written in the Write intermediate structure of one data node is performed in synchronization (Sync) with writing (update). It is said. By adopting such a configuration, it is possible to increase the response speed of Write for data that is not immediately accessed by the READ system for Write data.
0067When READ access is performed, it has already been converted to the data structure required for the READ access. Therefore, by processing the READ access using the converted data structure, the processing speed can be increased. be able to. Furthermore, depending on the type of READ access, it is possible to select an appropriate data structure and use the access destination node properly.
0068In addition, in FIG. 6 and the like, the number of types of data structures is set to three, A, B, and C, for the sake of simplicity of explanation, but the number of types of data structures is not limited to three. Of course, for example, any plurality of types having different characteristics may be used. Further, as an example of the data structure, three types of queue, column store, and row store have been illustrated, but it goes without saying that the data structure is not limited to such an example. For example Presence / absence of index in low store structure, -Difference in the type of column that created the index, -Low store format that stores updates in an additional structure, And so on.
0069In this way, by having an intermediate structure with write priority and asynchronously converting the structure, it is possible to avoid the bottleneck of structural conversion and maintain availability. In addition, by making it possible to control the data allocation node, data structure, and the trigger for applying asynchronous conversion (timer timeout time), the margin for various applications and load fluctuations is expanded.
0070Duplicate different data structures in the Sync method has a large overhead. Use a data structure such as a queue / log in the First-in, First-out (FIFO) method as a Write intermediate structure, and temporarily store the data in the intermediate structure. If it is reflected later, the conversion process will be more efficient and will have less impact on the access performance of the system.
0071By the way, in the configuration shown in FIG. 6, the trigger for asynchronously performing data conversion (the set value of the asynchronous timer (Async (timer)) in FIG. 6) is always optimal according to the data usage status. Is not always.
0072The setting value of the asynchronous timer in Fig. 6 is short, and frequent data structure conversion may adversely affect the write performance of the system. Conversely, if the asynchronous timer setting value (timeout time: update trigger information) in Fig. 6 is long and the frequency of data structure conversion is low, the system (analysis system) that uses the converted data structure is the latest. The data is not guaranteed to be analyzed, and there may be problems with the reliability of the analysis results.
0073That is, in the data node shown in FIG. 7, if the set value (timeout time) of the timer (Async (timer)) that defines the trigger for data structure conversion is relatively large, the data node will accumulate data in the Write intermediate structure. , It takes a long time to convert to the target data structure (column store format in Fig. 7). That is, in the data node, conversion to the target data structure and storage in the data storage unit are rarely performed, and data is exclusively stored in the Write intermediate structure. In this case, it is advantageous for the performance of the Write system. Further, since the data stored in the Write intermediate structure may be collectively converted into a data structure (for example, a column store format), the conversion process by the data structure conversion means (113 in FIG. 2) is also efficient.
0074However, in the data node, it takes a long time from receiving the data to converting the data into the target data structure, and a batch processing client (target data structure) that operates in batch processing or the like at a predetermined time or time zone. The data converted to is analyzed by batch processing), the old data whose data structure has been converted (old data) will be analyzed. When the latest or new data is required, the data stored in the Write intermediate structure of the data node (waiting for conversion of the data structure) is read and the data structure is converted to the column store format which is the target data structure ( (New data), the analysis will be performed after reflecting the difference between the old and new data in the column store format. In this case, the load on the client side increases.
0075On the other hand, if the set value (timeout time: update trigger information) of the asynchronous timer (Async (timer)) that defines the trigger for data structure conversion is relatively small in the data node, the data node will receive the received data. , It must be converted to the desired data structure little by little at short time intervals. Therefore, when the set value of the asynchronous timer (Async (timer)) is small, the write performance of the data node is disadvantageous as compared with the case where the set value of the asynchronous timer (Async (timer)) is large. On the other hand, when the setting value of the asynchronous timer (Async (timer)) is small, for example, a client (batch processing client) that analyzes data in batch processing can always refer to new data. Also, since the data whose data structure has been converted asynchronously is relatively recent data, the amount of data read from the Write intermediate structure is small even when the client refers to newer data, and the client side The load is also small.
0076In the configuration of FIG. 6, the trigger of data structure conversion by the asynchronous method in each data node depends on, for example, the method of data reference (Read access) from the client side.
0077<Write intermediate structure: embodiment> Therefore, in the present embodiment, as shown in FIG. 8, for example, the trigger for data structure conversion (update trigger information in FIG. 3A) is adjusted in relation to the frequency of access. If the access frequency (read access frequency) is below / above a predetermined threshold, increase / decrease the set value (timeout time) of the asynchronous timer (Async (timer)). That is, the value of the update trigger information (asynchronous timer: update trigger information in FIG. 3 (A)) of the data structure management information 921 (Fig. 2) is adjusted according to the access frequency.
0078When the load of the write system is larger / smaller than the load of the read system (analysis system), the set value (timeout time) of the asynchronous timer (Async (timer)) is increased / decreased. That is, when the frequency of Write access is higher than the frequency of Read access, the set value (timeout time) of the asynchronous timer is set large.
0079Alternatively, if the reference access (Read access) pattern is periodic (for example, when Read access is performed regularly) based on the access history information, the reference timing (Read access date / time, time zone, etc.) is adjusted. , Write The data accumulated in the intermediate structure may be converted into a target data structure and stored, and after the conversion, the set value (timeout time) of the asynchronous timer (Async (timer)) may be increased. Alternatively, the number of data structure conversions can be reduced by increasing the set value (timeout time) of the asynchronous timer (Async (timer)) because it is sufficient to make it in time for the next Read access that is performed periodically. Although not particularly limited, the data structure of the latest data may be set to be converted as much as possible before (immediately before) the next Read access is performed.
0080When the access history information is changed, for example, the value of the update trigger information (asynchronous timer timeout time) of the data structure management information 921 (Fig. 2) may be adjusted in synchronization (linkage) with this change.
0081According to this embodiment, the performance balance between the write system performance of online processing and the analysis system (read system) of batch processing is optimized only by adjusting the value of the update trigger information (asynchronous timer). Can be done.
0082In addition, in FIG. 8, the access frequency is shown in order to clarify the relationship with the change of the setting value of the asynchronous (Async) timer, and the configuration in which the access frequency information is stored and held in the data node is shown. However, the access frequency information of the data node may be provided outside the data node. Alternatively, the access frequency information of the data nodes may be stored and managed in a common storage for a plurality of data nodes. In addition, in the data node, the access history (log) is taken, the access frequency is calculated based on the access history information, and the setting value (update trigger information) of the asynchronous (Async) timer is changed based on the access frequency. You may do it. Alternatively, instead of the access frequency (the number of applications for access in a unit period), the set value (update trigger information) of the asynchronous (Async) timer may be changed by using the access tendency, the access pattern indicating the characteristic, or the like. Good.
0083<Change judgment means> FIG. 9 is a diagram showing an example of a configuration for adjusting the update trigger information of the data structure management information 921. As shown in FIG. 9, the change determination means (change determination unit) 72 for determining whether or not to change the update trigger information of the data structure management information 921 is provided based on the access information of the access history recording unit 71. ..
0084As described with reference to FIG. 2, the access receiving means 111 of each data node records the received access request in the access history recording unit 71. The access history recording unit 71 records the access request (including the table identifier in FIG. 3A, the replica identification value of the data node, etc.) in association with the time information (date and time information) at the time of receiving the access request.
0085Although the access history recording unit 71 is provided for each data node, it may be provided with one for a data node group composed of a plurality of data nodes, or one for the entire system. Alternatively, each data node may be provided with an access history recording unit 71, and a mechanism for aggregating the access frequency information individually collected by each data node may be provided by an arbitrary method.
0086The change determination means (change determination unit) 72 uses the access history information stored in the access history recording unit 71 to determine the frequency of access (threshold value) within a period of the most recent predetermined length, for example. Whether or not to change the update trigger information related to the corresponding data node according to the comparison result with<u style="single">Decision</u>You may try to do it. Alternatively, the access frequency within the period of the most recent predetermined length is calculated, and the magnitude (threshold value) of the fluctuation from the value of the access frequency information in the period of the previous predetermined length is calculated. Whether or not to change the update trigger information related to the corresponding data node according to the comparison result)<u style="single">Decision</u>You may try to do it.
0087The change determination means 72 sets the asynchronous timer for the structural information changing means 91 when it is necessary to change the update trigger information (setting value of the timeout time of the asynchronous timer) for asynchronous conversion in the related data node. Issue a value change request. The change request from the change determination means 72 includes the replica identifier corresponding to the data node, the table identifier information, and the node information of the data node. Further, the change request from the change determination means 72 does not change (change value = 0), increments / decrements by a predetermined unit, or increases or decreases by a multiple of a predetermined unit with respect to the current asynchronous timer setting value. May include the instruction. Alternatively, the change determining means 72 derives the changed value of the set value of the asynchronous timer, sets this changed value in the change request, and the structural information changing means 91 replaces the current set value of the asynchronous timer with the changed value. It may be configured. The relationship between the table identifier information, the replica identifier, and the data node information (placement node number) is defined in the data placement specific information 922, and the structural information changing means 91 responds to the change request from the change determining means 72. Then, from the data node information (ID), replica identifier, and table identifier information, the corresponding table identifier information and replica identifier update trigger information in the data structure management information 921 are changed.
0088In FIG. 9, the change determination means 72 is provided separately from the data management / processing means 11 of the data node 1, but the change determination means 72 is implemented in the data management / processing means of each data node. Then, the access history recording unit 71 may hold the access frequency information calculated by the change determination means 72.
0089The access frequency information is not necessarily limited to the number of times a Read access request is generated / the number of times a Write access request is generated per unit period, and is not necessarily limited to, for example, the pattern of occurrence of Read and Write access requests (Read, Write). If access occurs at a fixed time or the like, the information in the timetable) may be used.
0090<Client access flow> FIG. 10 is a flowchart for explaining the operation of the client function realization means 61 of FIG. 1 in which the client function realization means 61 issues an instruction to the update destination data node and waits for the data node. The client access flow will be described with reference to FIG.
0091The client function realizing means 61 acquires the information of the structural information holding unit 92 by accessing the master data (master file) or a cache at an arbitrary location (a cache memory storing a copy of a part of the master data). (Step S101 in FIG. 10).
0092Next, the client function implementation means 61 identifies whether the instruction content issued by the client is a WRITE process or a reference process (Read) (step S102).
0093This can be specified by specifying it with the command of the issue instruction or by analyzing the execution code of the instruction. For example, in the case of a storage system that processes SQL -If it is an INSERT instruction (SQL instruction that adds a record to a table), WRITE processing, -If it is a SELECT instruction (SQL instruction that refers to and retrieves records from a table), reference processing, Is.
0094Alternatively, the client function implementation means 61 may be used to explicitly specify the instruction when the instruction is called (such an API (Application Program Interface) is prepared).
0095If the result of step S102 is WRITE processing, the process proceeds to step S103 and subsequent steps.
0096In the case of WRITE processing, the client function realization means 61 identifies the node that needs to be updated by using the information of the data arrangement identification information 922.
0097The client function implementation means 61 issues an instruction execution request (update request) to the specified data node (step S103).
0098The client function implementation means 61 waits for a response notification from the data node to which the update request is issued, and confirms that the update request is held in each data node (step S104).
0099If the result of step S102 is reference processing, the process proceeds to step S105.
0100In step S105, the client function realizing means 61 identifies (recognizes) the characteristics of the processing content (step S105).
0101Next, the client function realizing means 61 selects a data node to be accessed based on the specified processing characteristics and other system conditions, and performs a process of issuing an instruction request (step S106).
0102The client function implementation means 61 then receives the access processing result from the data node (step S107).
0103Hereinafter, the description of the processing in steps S105 and S106 will be supplemented. The client function realizing means 61 can know the type of data structure in which the data to be accessed is held from the information stored in the data structure management information 921. For example, in the case of FIG. 3A, when accessing the WORKERS table, the replica identifiers 0 and 1 are the data structure B, and the replica identifier 2 is the data structure C. In the access frequency information, the access to the WORKERS table is recorded in association with the replica identifier of the data node.
0104Then, the client function realizing means 61 determines which data structure is suitable for the data access performed to the data node, and selects the suitable data structure. More specifically, for example, in the client function realization means 61, when the access is an access that parses the SQL statement that is the access request and the table identifier is the sum of a certain column in the table of "WORKERS", the data structure C ( Column store) is selected. When the SQL statement is an access to retrieve a specific record, the client function implementation means 61 determines that the data structure B (low store) is suitable.
0105In the case of an instruction to retrieve a specific record, the client function realization means 61 may select either of the replica identifiers 0 and 1. It is desirable to use the replica identifier 1 in which the update trigger information is set to a large value, "when it is not always necessary to perform processing with the latest data".
0106The identification of this "when there is no need to process with the latest data" depends on the application context. Therefore, the data structure to be used and the information for specifying the freshness of the required data (newness of the data) may be explicitly specified in the instruction passed to the client function realizing means 61.
0107The client function implementation means 61 calculates the data node to be accessed after specifying the replica identifier (data structure) to be accessed. At this time, the selection of the access node may be changed according to the situation of the distributed storage system. For example, when a certain table is stored in the data nodes 1 and 2 as the same data structure B and the access load of the data node 1 is large, the client function realizing means 61 selects the data node 2. You may change to the operation.
0108Further, as another data structure C, when the access load of the data node 3 is smaller than that of the data nodes 1 and 2 when stored in the data node 3, the access content to be processed is the data structure B. Even if it is more suitable, the client function realization means 61 may issue an access request to the data node 3 (data structure C).
0109The client function realizing means 61 issues an access request to the data node calculated and selected in this way (S106), and receives the access processing result from the data node (S107).
0110<Data node operation> FIG. 11 is a flowchart illustrating the access process in the data node of FIG. The operation of the data node will be described in detail with reference to FIGS. 11 and 2.
0111First, the access receiving means 111 of the data management / processing means 11 of the data node receives the access processing request (step S201 in FIG. 11).
0112Next, the access receiving means 111 of the data management / processing means 11 of the data node determines whether the content of the received processing request is a Write process or a Read (reference) process (step S202).
0113If the result of step S202 is WRITE processing, the access processing means 112 of the data management / processing means 11 of the data node acquires the information of the data structure management information 921 in the structure information holding unit 92 (step S203). The information acquisition of the data structure management information 921 may access the master data or access the cache data (data in the cache memory storing a partial copy of the master data) at an arbitrary location. Alternatively, the client function realizing means 61 of FIG. 1 adds information (access to master data or cache data) to the request issued to the data node, and the access processing means 112 uses the information. You may want to access it.
0114Next, the access processing means 112 determines from the information of the data structure management information 921 whether or not the update trigger of the processing for the data node is 0 (zero) (step S204).
0115As a result of step S204, when the update trigger is "0", the access processing means 112 directly updates the data structure specified in the structural information of the structural information holding unit 92 (step S205). That is, the update data is converted into the specified data structure and stored in the corresponding structure-specific data storage unit 12X (X = 1, 2, 3).
0116If the update trigger is not "0", the access processing means 112 stores the update data in the Write intermediate structure (structure-specific data storage unit 121) (step S206).
0117In the case of steps S205 and 206, after the processing is completed, the access receiving means 111 responds to the requesting client function realizing means 61 with the processing completion notification (step S207).
0118If the result of step S202 is data reference processing, the reference processing is executed (step S208).
0119The execution method of the Read (reference) process is not particularly limited, but the following three types of methods can be typically mentioned.
0120(1) The first method uses the data in the data storage section of the data structure specified in the data structure management information 921 for processing. This has the best performance, but if the update trigger time (cycle) is long, there is a possibility that the data of the Write intermediate structure is not reflected in the reference processing. This can lead to data inconsistencies. However, if the application developer is aware of it in advance and uses it, or if it is known that data reading does not occur within the update trigger after writing, or if new data access is required, the update trigger will occur. If you decide to access the "0" replica identifier data, there is no problem.
0121(2) The second method is to wait for the application of the conversion process to be performed separately before processing. This is easy to implement, but the response performance deteriorates. For applications that do not require responsiveness, there is no problem.
0122(3) The third method reads and processes both the data structure specified in the data structure management information 921 and the data held in the Write intermediate structure. In this case, the latest data can always be responded, but the performance is worse than the first method.
0123Any of the above-mentioned first to third methods may be taken. Further, it is possible to specify the method to be executed in the processing instruction issued from the client function realization means 61, which realizes a plurality of types and describes it as a system configuration file.
0124<Data structure conversion operation of data structure conversion means> FIG. 12 is a flowchart showing the operation of the data conversion process in the data structure conversion means 113 of FIG. The data conversion process will be described with reference to FIGS. 12 and 2.
0125The data structure conversion means 113 periodically waits for a call due to the occurrence of a timeout in a timer (not shown) in the data node in order to determine whether or not conversion processing is necessary (step S301 in FIG. 12). Note that this timer may be provided in the data structure conversion means 113 as a dedicated timer. The timer timeout time corresponds to the update trigger information (sec) setting value in FIG. 3 (A) (Aync (timer) timeout time in FIG. 6).
0126Next, the structural information (data information) of the structural information holding unit 92 is acquired (step S302), and it is determined whether or not there is a data structure that needs to be converted (step S303). For example, when the timer makes a determination every 10 seconds, the data structure whose update trigger is 20 seconds executes the conversion process every 20 seconds, so that the conversion process does not have to be performed at the time of 10 seconds. If no conversion process is required, the process returns to waiting for a timer call (waiting until it is called due to a timeout occurring in the timer) (step S301).
0127On the other hand, when conversion processing is required, the update processing content for the data to be converted is read from the intermediate data structure for update (step S304), and the data storage unit 12X (X = 1 to 3) for each structure of the conversion destination is read. Perform the process to reflect the update information (step S305).
0128<Write sequence 1> FIG. 13 is a diagram showing a sequence of Write processing (processing accompanied by updating data).
0129The client function realization means 61 (client computer) of the client node 6 acquires the information of the data arrangement specific information 922 (Fig. 2) held in the structural information holding unit 92 of the structural information management means 9 (or at an arbitrary location). Get information from cache memory).
0130The client computer uses the acquired information to issue a Write access instruction to the data node (data node 1 with replica identifier 0) to which the data to be written is placed.
0131The access receiving means 111 of the data node 1 receives the write access request and transfers the write access to the data nodes 2 and 3 specified by the replica identifiers 1 and 2. As a method of identifying the data nodes of the replica identifiers 1 and 2, the data node 1 may access the structural information holding unit 92 (or an appropriate cache), or the Write access instruction issued by the client function realizing means 61 may be used. All or part of the data structure management information 921 may be passed together.
0132The access processing means 112 of each data node processes the received Write access request.
0133The access processing means 112 refers to the information of the data structure management information 921 and executes the write process.
0134If the value of the update trigger information is larger than "0", the Write processing content is stored in the structural data storage unit 121 of the data structure A.
0135When the value of the update trigger information is "0", it is stored in the structure-specific data storage unit 12X of the data structure specified in the data structure management information 921.
0136After the write process is completed, the access processing means 112 issues a completion notification to the access receiving means 111 and returns a completion response to the client computer.
0137The replica destination data node (2, 3) returns a write completion response to the access receiving means 111 of the replica source data node 1.
0138The access receiving means 111 waits for the completion notification from the access processing means 112 of the data node 1 and the completion notification of the data nodes 2 and 3 of each replica destination, and after receiving all of them, returns a response to the client computer.
0139The data structure conversion means 113 (see FIG. 2) of the data node 1 converts the data stored in the Write intermediate structure (structure-specific data storage unit 121 (data structure A)) into the structure-specific data according to the timeout of the asynchronous timer. Converts to storage unit 12X (final storage destination data structure specified in data structure management information 921) and stores it. Similarly, data nodes 2 and 3 also perform conversion to the target data structure according to the timeout of the asynchronous timer.
0140<Write sequence 2> In the example of FIG. 13, the data node 1 transfers the write request to the replica destination data nodes 2 and 3, but as shown in FIG. 14, the client computer uses the storage destination data node. A Write request may be issued for all of the above.
0141In the example of FIG. 14, the client computer waits for the Write access request, which is different from that of FIG. In the example of FIG. 14, the client computer issues a write request to the storage destination data nodes 0, 1, and 2, respectively, and receives a completion response from the storage destination data nodes 0, 1, and 2, respectively.
0142<Modification example> FIG. 15 is a diagram illustrating a modified example of the configuration of FIG. Referring to FIG. 15, the column store format data node 3 of FIG. 8 is composed of two data nodes 3A and 3B, and one data node 3A from the write intermediate structure to the column store (Column). When converting to the Store) format data structure, the analysis client (Client) has already converted the data of the other data node 3B (the data before conversion stored in the Write intermediate structure and the column store format). Analyze by referring to the data). The asynchronous timer setting on data nodes 3A and 3B is 20 seconds (Aync (20 seconds)), but data structure conversion on data node 3B is 10 seconds more than data structure conversion on data node 3A. Running late. For example, in the data node 3A, the data structure is converted in the time interval of 0 to 20 seconds, and the data is analyzed by the client (Client) that performs Read access in the following time interval of 20 to 40 seconds. In the data node 3B, the data is analyzed by the client that performs Read access in the time interval of 10 seconds to 30 seconds, and the data structure is converted in the following time interval of 30 to 50 seconds. Therefore, for example, at 15 seconds, which is between 10 seconds and 20 seconds, data structure conversion is performed at data node 3A, and data analysis is performed at data node 3B. The asynchronous timer settings in the data nodes 3A and 3B are set based on the access history information (access frequency) of the data nodes 3A and 3B.
0143<Another variant> FIG. 16 shows an example in which an ETL (Extract / Transform / Load) is arranged between an online processing (online processing system that performs Write processing) and an analysis system (data warehouse) that is performed in batch processing or the like.
0144A data warehouse system includes a large-scale database for information analysis and decision making by extracting and reconstructing data (for example, transaction data) from a core system. It is necessary to migrate data from the database of the core system to the data warehouse database, and this process is called ETL (Extract / Transform / Load). In addition, "Extract" extracts data from the information source of the department, "Transform" transforms and processes the extracted data as needed in the business, and "Load" converts it to the final target (that is, data warehouse). Represents loading processed data. In FIG. 16, the above embodiment is applied to ETL data conversion. That is, the asynchronous data conversion by the ETL in FIG. 16 corresponds to the data structure conversion by the data structure conversion means 113 in FIG.
0145In the example of FIG. 16, the ETL asynchronously converts the row store format data (replica data) of the active system (online processing) to the column store (Column-Store) format for the analysis system (data warehouse). Converted with (Asynch: Asynchronous). In the present embodiment, the bottleneck of data structure conversion is eliminated and the storage utilization efficiency is improved by adjusting the timer for asynchronous conversion in ETL based on the access frequency based on the access history information (access frequency information). Can be enhanced.
0146Each disclosure of the above patent documents shall be incorporated into this document by citation. Within the framework of the entire disclosure (including the scope of claims) of the present invention, it is possible to change or adjust the embodiments or examples based on the basic technical idea thereof. Further, various combinations or selections of various disclosure elements (including each element of each claim, each element of each embodiment, each element of each drawing, etc.) are possible within the scope of the claims of the present invention. .. That is, it goes without saying that the present invention includes all disclosure including claims, and various modifications and modifications that can be made by those skilled in the art in accordance with the technical idea.
01471-4 data nodes 5 network 6 client node 9 Structural information management means (structural information management device) 11, 21, 31, 41 Data management / processing means (data management / processing department) 12, 22, 32, 42 data storage 61 Client function realization means (Client function realization department) 71 Access history recording section 72 Change judgment means (change judgment unit) 91 Structural information changing means (Structural information changing part) 92 Structural information holding unit 111 Access reception means (access reception department) 112 Access processing means (access processing unit) 113 Data structure conversion means (data structure conversion unit) 121, 122, 123, 12X Structural data storage 611 Data access means (data access section) 612 Structural information cache holder 921 Data structure management information 922 Data placement specific information
16 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8 Sheet 9 Sheet 10 Sheet 11 Sheet 12 Sheet 13 Sheet 14 Sheet 15 Sheet 16
Every citation, both ways
| Document | Relation | Office |
|---|---|---|
| JP2010128752A | Cites | Japan |
| JP2002197781A | Cites | Japan |
| JP2007317017A | Cites | Japan |
| US20080154914A1 | Cites | United States of America |
| WO2010101189A1 | Cites | World Intellectual Property Organization (WIPO) |
| US20110320750A1 | Cites | United States of America |
| EP2405360A1 | Cites | European Patent Office (EPO) |
| CN102341791A | Cites | China |
5 members in 3 offices
Priority claims3
| Document | Office | Kind | Date |
|---|---|---|---|
| 2011169588 | Japan | – | |
| 2011169588 | Japan | A | |
| 2012069499 | Japan | W |
Members5
| Document | Office | Kind | |
|---|---|---|---|
| WO2013018808A1 | World Intellectual Property Organization (WIPO) | A1 | |
| US2014173035A1 | United States of America | A1 | |
| JPWO2013018808A1 | Japan | A1 | |
| JP6044539B2This record | Japan | B2 | |
| US9609060B2 | United States of America | B2 |
8 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Cancellation because of no payment of annual feesLAPS | LAPS | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Request for written amendment filedJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Notification of reasons for refusalJAPANESE INTERMEDIATE CODE: A131A131 | A131 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 6044539
- Application
- 2013526936
Titles2
- Japanese
- 分散ストレージシステムおよび方法
- English
- Distributed storage system and method
Classification
- CPC, 9
- H04L67/1097
- H04L67/1095
- G06F3/061
- G06F3/065
- G06F3/067
- G06F3/0617
- G06F11/2094
- G06F16/184
- H04L67/568
- IPC, 1
- G06F12 00
