Data processing program, server apparatus, and data processing method
Abstract
Problem to be solved.To dynamically determine the number of replicas of each data and to efficiently execute data processing in each server device.
Solution.In a server device 100, a number of replicas indicating how many server devices 100 are to be replicated is set for each data, and based on the number of replicas, which server device 100 to allocate data to is selected. .. Then, each server device 100 redetermines the server device 100 to which the data is arranged every time the number of replicas changes dynamically. The number of replicas increases, and the server device 100 newly added as a data allocation destination newly allocates the corresponding data. On the other hand, the number of replicas decreases, and the server device 100 excluded from the data allocation destination deletes the corresponding data. [Selection diagram] Fig. 1

Term
Projected expiry 16 December 2028.
- Priority and filed
- Published
- Today
- Projected expiry
5 claims: 3 independent, 2 dependent
- 1相互に通信可能なコンピュータ群を構成するコンピュータを、 任意のデータについての処理要求が入力されると、当該任意のデータに設定されている複製数を取得する取得手段、 前記コンピュータ群の中から、前記任意のデータの配置先となるコンピュータを、所定のアルゴリズムを用いて前記複製数分選択する選択手段、 前記取得手段によって取得された前記任意のデータの複製数を前記コンピュータ群すべてに送信する複製数送信手段、 前記選択手段によって選択された前記複製数分の各コンピュータに、前記処理要求を送信する処理要求送信手段、 自装置、または、他のコンピュータから送信された処理要求を受け付けた場合、当該処理要求に応じた処理を実行する実行手段、 任意のタイミングごとに、前記実行手段によって実行された前記任意のデータへの処理要求を参照して前記任意のデータの複製数を決定する決定手段、として機能させ、 前記複製数送信手段は、前記決定手段によって現在設定されている複製数とは異なる複製数が決定された場合に、当該決定された複製数を前記コンピュータ群すべてに送信し、 前記選択手段は、前記複製数送信手段によって前記決定された複製数が送信されてきた場合、あらたに、前記コンピュータ群の中から、前記任意のデータを配置するコンピュータを、所定のアルゴリズムに応じて前記決定された複製数分選択し、 前記実行手段は、自装置が、前記選択手段によって、あらたに前記任意のデータを配置するコンピュータに選択された場合に前記任意のデータを書き込み、前記選択手段によって、あらたに前記データを配置するコンピュータに選択されなくなった場合に前記データを削除することを特徴とするデータ処理プログラム。
- 2前記決定手段は、前記コンピュータ群の総数を前記任意のデータの複製数とした場合の前記任意のデータに対する処理時間の平均値が所定値以上の場合、前記コンピュータ群の総数を前記任意のデータの複製数に決定することを特徴とする請求項1に記載のデータ処理プログラム。
- 3前記決定手段は、前記任意のデータの複製数を1とした場合の前記任意のデータへの処理時間の平均値が所定値未満であった場合、前記任意のデータの複製数として設定可能な最小値を複製数に決定することを特徴とする請求項1または2に記載のデータ処理プログラム。
- 4相互に通信可能なサーバ装置群を構成するサーバ装置であって、 任意のデータについての処理要求が入力されると、当該任意のデータに設定されている複製数を取得する取得手段と、 前記サーバ装置群の中から、前記任意のデータの配置先となるサーバ装置を、所定のアルゴリズムを用いて前記複製数分選択する選択手段と、 前記取得手段によって取得された前記任意のデータの複製数を前記サーバ装置群すべてに送信する複製数送信手段と、 前記選択手段によって選択された前記複製数分の各サーバ装置に、前記処理要求を送信する処理要求送信手段と、 自装置、または、他のサーバ装置から送信された処理要求を受け付けた場合、当該処理要求に応じた処理を実行する実行手段と、 任意のタイミングごとに、前記実行手段によって実行された前記任意のデータへの処理要求を参照して前記任意のデータの複製数を決定する決定手段と、を備え、 前記複製数送信手段は、前記決定手段によって現在設定されている複製数とは異なる複製数が決定された場合に、当該決定された複製数を前記サーバ装置群すべてに送信し、 前記選択手段は、前記複製数送信手段によって前記決定された複製数が送信されてきた場合、あらたに、前記サーバ装置群の中から、前記任意のデータを配置するサーバ装置を、所定のアルゴリズムに応じて前記決定された複製数分選択し、 前記実行手段は、自装置が、前記選択手段によって、あらたに前記任意のデータを配置するサーバ装置に選択された場合に前記任意のデータを書き込み、前記選択手段によって、あらたに前記データを配置するサーバ装置に選択されなくなった場合に前記データを削除することを特徴とするサーバ装置。
- 5相互に通信可能なコンピュータ群を構成するコンピュータが、 任意のデータについての処理要求が入力されると、当該任意のデータに設定されている複製数を取得する取得工程と、 前記コンピュータ群の中から、前記任意のデータの配置先となるコンピュータを、所定のアルゴリズムを用いて前記複製数分選択する選択工程と、 前記取得工程によって取得された前記任意のデータの複製数を前記コンピュータ群すべてに送信する複製数送信工程と、 前記選択工程によって選択された前記複製数分の各コンピュータに、前記処理要求を送信する処理要求送信工程と、 自装置、または、他のコンピュータから送信された処理要求を受け付けた場合、当該処理要求に応じた処理を実行する実行工程と、 任意のタイミングごとに、前記実行工程によって実行された前記任意のデータへの処理要求を参照して前記任意のデータの複製数を決定する決定工程と、を実行し、 さらに、 前記複製数送信工程では、前記決定工程によって現在設定されている複製数とは異なる複製数が決定された場合に、当該決定された複製数を前記コンピュータ群すべてに送信し、 前記選択工程では、前記複製数送信工程によって前記決定された複製数が送信されてきた場合、あらたに、前記コンピュータ群の中から、前記任意のデータを配置するコンピュータを、所定のアルゴリズムに応じて前記決定された複製数分選択し、 前記実行工程では、自装置が、前記選択工程によって、あらたに前記任意のデータを配置するコンピュータに選択された場合に前記任意のデータを書き込み、前記選択工程によって、あらたに前記データを配置するコンピュータに選択されなくなった場合に前記データを削除することを特徴とするデータ処理方法。
Independent claims5
142 paragraphs, as filed
The present invention relates to a data processing program, a server device, and a data processing method in each computer constituting a group of computers capable of communicating with each other.
In recent years, in database systems that can be accessed via network systems, as the amount of data to be placed increases and the number of accesses to data increases, data is stored in multiple servers and other devices such as storage (hereinafter, these are collectively referred to as "servers". The number of configurations that are distributed and managed in (referred to as "devices") is increasing. When arranging data in a distributed manner in this way, it is possible to create a replica of one data and allocate it to each of multiple servers in order to distribute the load and improve availability. There are many.
The total number of replicas including the original data generated at this time is called the number of replicas. Normally, the number of replicas is determined according to the access pattern to the data. For example, if all the data is only referenced, then all the data is copied to all the server devices where the data is located. That is, the number of replicas = the total number of server devices. Such dispersion has the following advantages.
-Since data is distributed to all server devices, the processing load due to reference requests can be easily distributed by preparing a load balancer with a simple configuration that distributes reference requests from clients to one of the server devices. be able to. -Even if there are reference requests for the same data from multiple clients, the reference requests can be distributed by the number of servers. -High availability because data will not be lost unless all server devices go down.
On the other hand, when not only data reference but also data update is performed, if the data is copied to all server devices as described above, every time the update occurs, the data is updated to all server devices. Must reflect the update of. As a result, copy processing occurs frequently in each server device, resulting in poor processing efficiency. In addition, if the number of replicas is set to 1 (no copy) in order to maximize the processing efficiency of each server device even when updates are performed, the load will be increased when reference requests for the same data are concentrated. It may not be possible to distribute and the response time to the reference request may become long. Furthermore, since a spare server device in which the same data is placed is not prepared in addition to the server device in which the data is placed, there is a risk that the data will be lost if the server device goes down, and the availability will be significantly reduced. ..
Therefore, conventionally, the number of replicas is often set to 2 or more and less than the total number of server devices. In practice, the number of replicas is determined in consideration of the data access pattern (reference / update rate, frequency) and the availability required by the database administrator or user. This enables efficient data access for both reference and update.
<patcit num="1"><text>Japanese Unexamined Patent Publication No. 2-231676</text></patcit>
<p> However, as described above, when the access pattern to each data arranged in the server device used as a data distributed database system differs greatly depending on the data, the optimum number of replicas common to all the data is determined. That was difficult. For example, if the number of replicas (common to all data) is set high, the reference efficiency for data with many references will increase, but for data with many updates, the data for the number of replicas must be updated, and when updating Processing efficiency is reduced.</p><p> As a countermeasure against the above-mentioned case, a technique is also provided in which data with many references and data with many updates are separated, the number of replicas is set for each data, and the data is arranged in each server device. However, if the data cannot be divided due to differences in access patterns, such as when a method of classifying data according to the data owner is applied, only an average value can be set as the number of replicas.</p><p> Further, even if data can be divided based on an access pattern at a certain point in time and different numbers of replicas can be set, it is difficult to always apply the optimum number of replicas if the access pattern changes frequently. For example, in the case of a service such as a blog where the update process to the database is mainly performed by a general user other than the server device administrator, the update frequency of the user's blog article may change due to the convenience of the user, or some news reports may occur. , The frequency of referencing a user's page suddenly increases due to the occurrence of an event or the like. Therefore, if the number of replicas remains fixed, there is a problem that the efficiency of data access is poor and it becomes difficult for the referrer to refer comfortably.</p><p> Therefore, in a data-distributed database system, the reference frequency of data is monitored, and the comparison result with a preset reference value (for example, when the reference value is exceeded or is expected to be exceeded), a part of the data. Alternatively, a technique of copying the entire data and sending it to another site has also been proposed (see, for example, Patent Document 1 above). According to this technique, the reference efficiency can be improved by copying the replica to a site with a high reference frequency. In addition, since the copy is dynamically created, it is possible to effectively respond to changes in the tendency of database processing.</p><p> However, since the above technology only describes copying data, if you continue to use this technology in an environment where data access changes drastically, all data will eventually be copied to all sites unless you limit the data copy destination. Will end up. Therefore, when the reference frequency of a certain data decreases and the update frequency increases, the effect of load distribution of the reference is small, and conversely, the update is inefficient due to the large number of copies.</p><p> In order to solve the problems caused by the above-mentioned conventional techniques, a data processing program, a server device, and data capable of dynamically determining the number of replicas of each data and efficiently executing data processing in each server device. The purpose is to provide a processing method.</p>
<p> In order to solve the above-mentioned problems and achieve the purpose, when a processing request for arbitrary data is input in a computer constituting a group of computers capable of communicating with each other, the number of duplicates set for the arbitrary data is set. The process of selecting the computer to which the arbitrary data is to be placed from the computer group for the number of duplicates using a predetermined algorithm, and the number of duplicates of the acquired arbitrary data. When the process of transmitting the processing request to all the computers, the process of transmitting the processing request to each of the selected computers for the number of duplicates, and the processing request transmitted from the own device or another computer are received. , A process for executing a process according to the process request, a process for determining the number of duplicates of the arbitrary data by referring to the executed process request for the arbitrary data at an arbitrary timing, and a current setting. When the number of duplicates different from the number of duplicates is determined, the process of transmitting the determined number of duplicates to all the computers and the case where the determined number of duplicates is transmitted, A process of newly selecting a computer for arranging the arbitrary data from the computer group for the number of duplicates determined according to a predetermined algorithm, and the own device newly arranging the arbitrary data. It is a requirement to include a process of writing the arbitrary data when selected by the computer to be used, and a process of deleting the data when the computer is no longer selected by the computer on which the data is newly arranged.</p><p> According to this data processing program, server device, and data processing method, even if the number of replicas changes dynamically for each data by monitoring not only the reference frequency of the data but also the update frequency at the same time, each data Select the location of each data based on the set number of replicas. Therefore, the optimum number of replicas can be set individually for each data, and the number of computers to which data is placed is dynamically changed according to the frequency of processing requests without being affected by the number of replicas of other data. be able to.</p>
<p> According to this data processing program, the server device, and the data processing method, it is possible to dynamically determine the number of replicas of each data and to efficiently execute the data processing in each server device.</p>
Preferred embodiments of the data processing program, server apparatus and data processing method will be described in detail below with reference to the accompanying drawings. In this data processing program, server device, and data processing method, the number of replicas indicating how many server devices are to be replicated is set for each data, and based on this number of replicas, which server device to allocate the data to is determined. select. Then, each server device reselects the server device to which the data is arranged every time the number of replicas changes dynamically. The number of replicas increases, and the server device newly added as the data allocation destination newly allocates the corresponding data. On the other hand, the number of replicas decreases, and the server device excluded from the data allocation destination deletes the corresponding data.
That is, it is possible to dynamically increase or decrease the number of server devices in which data is arranged according to a change in the number of replicas. Therefore, as in the past, reference requests from clients are concentrated and the response time becomes long, and the number of server devices on which the same data is placed is large even though the update frequency is high. It is possible to solve the problem that the server device is busy with the update process and the processing efficiency is lowered. Hereinafter, the best mode for realizing the above-mentioned data processing will be specifically described.
(Overview of data processing) First, an outline of data processing according to the present embodiment will be described. FIG. 1 is an explanatory diagram showing a system configuration of the server device according to the present embodiment. As shown in FIG. 1, in the present embodiment, a data distribution system 200 is realized in which data is distributed and arranged in 100 groups of server devices having the same configuration. The data distributed by the data distribution system 200 is not particularly limited to weblog articles and shared databases. Therefore, any system that handles data that is expected to be referenced or updated by a plurality of users can be adapted to various uses.
A load balancer 110 that receives a request 120 representing a processing request from the user is connected to the data distribution system 200. The load balancer 110 allocates the request 120 received from the outside to any of the server devices 100-1 to the server device 100-n prepared as the data distribution system 200.
When the request 120 is allocated from the load balancer 110, the server device 100 executes the application (application according to the content of the request 120) in the application execution unit 101 to perform the data processing requested as the request 120. Specifically, the data processing executed in the server device 100 includes two types, an update process to the target data and a reference process, as shown below.
<Data update process> It is a process that changes the content of the target data specified by request 120, and is classified into three types: a new write process, an existing data update process, and an existing data deletion process. <Data reference processing> This is a process in which the content of the target data specified by request 120 does not change, and the existing data reading process corresponds to this.
The request 120 described in the present embodiment is composed of information for identifying which of the above data processing requests is requested and information regarding the target data. For example, the configuration of the request 120 at the time of data update and the request 120 at the time of data reference is as follows.
Request for update Request example 1: New write Access type: Update (create new file) Data information: New data (new file name and contents to be written) Request example 2: Update existing data Access type: Update (file overwrite) Data information: Data to be updated (file name to be overwritten and content to be overwritten) Request example 3: Delete existing data Access type: Update (Delete File) Data information: Data to be deleted (file name to be deleted) Request at the time of reference Request Example 4: Read existing data Access type: Browse (read file contents) Data information: Data to be read (file name to be read)
Then, in the case of the present embodiment, each server device 100 constituting the data distribution system 200 can further realize the following functions by including the cooperation processing unit 102.
1) The server device 100 that receives the request 120 for accessing (reference / updating) the target data from the outside such as a client sends the received request 120 to the server device 100 in which the data to be processed is arranged. .. At this time, if the own device is the target data placement location, the data is transmitted to the own device, and if the other server device 100 is the target data placement location, the data is also transmitted to the server device 100. When the request 120 is an update process, the content of the data newly written during the request or the content of the update data is included.
2) The server device 100 in which each data to be processed is arranged monitors the data access pattern (reference / update ratio, frequency), and periodically determines the number of replicas for each data.
3) When the number of replicas of data is changed, the server device 100 in which any of the data is arranged transmits the information to the other server device.
4) When the server device 100 in which any of the data is arranged receives the information regarding the change in the number of replicas of a certain data, the server device 100 specifies the allocation destination of the data from the changed number of replicas. When the number of replicas increases, the source server to send the missing replicas is determined using a predetermined rule, and if the own device is the source server, the replica of the data is copied to a new location. When the number of replicas decreases, if the own device deviates from the placement destination, the placed data is deleted.
The above 1) is a particular feature of the function of the server device 100 according to the present embodiment. As shown in FIG. 1, the processing request 120 for the data arranged in the data distribution system 200 (in the case of new writing, the data to be newly arranged) configures the data distribution system 200 via the load balancer 110. It is sent to one of the server devices 100. At this time, the server device 100 to which the request 120 is assigned selects the server device 100 in which the target data is arranged based on the number of replicas set in the target data of the request 120.
When the data distribution system 200 is configured by n server devices 100 as shown in Fig. 1, the target data with the number of replicas: 3 is 3 server devices 100 (for example, server devices 100-1,100-2,100-n). ) Is located. Therefore, when any of the server devices 100 receives the update request to the target data as request 120, the update request is transmitted to all three server devices 100-1, 100-2, 100-n in which the target data is arranged.
On the other hand, when any of the server devices 100 receives the reference request for the target data as the request 120, the reference request is transmitted to any of the three server devices 100-1, 100-2, 100-n. The server device 100 that has received the reference request refers to the arranged target data, and returns the reference result to the server device 100 that has received the request 120. In this way, in the data distribution system 200, two-way communication for processing the request is performed between the server devices 100. These two-way communication is directly performed via the cooperation processing unit 102 (described later) included in the server device 100.
As described above, in the present embodiment, each target data is distributed and arranged in any one of a plurality of server devices 100 prepared based on the number of replicas. At this time, the number of distributed units = the number of replicas. In addition, the number of replicas is calculated periodically while monitoring the data access status from the client. Then, the data placement destination is determined again based on the newly calculated number of replicas. Therefore, even in an environment where the access pattern changes frequently, the optimum number of replicas is always applied according to the change, and the performance as the data distribution system 200, the user who sends the request 120 to the target data, and the data distribution system 200 The convenience for the administrator can be kept constant.
In the conventional technology, when it is desired to avoid manually changing the number of replicas after the start of operation, it is necessary to analyze the access pattern to each data in detail in advance to determine the number of replicas. In the present embodiment, since each server device 100 can dynamically change the data arrangement contents according to the number of replicas, the number of replicas can be dynamically changed even when each server device 100 is in operation. can do. Therefore, even if the number of replicas is determined without detailed analysis in advance, the number of replicas is changed while analyzing the actual data, so that the optimum number of replicas can be finally applied.
In addition, when the number of replicas is determined by prior analysis as in the conventional technique, the data in actual operation cannot be used, and the predicted value may be obtained. In many cases, the predicted value is far from the actual situation, and as a result, even if a detailed analysis is performed, the analysis result does not maintain the efficiency of the server device 100 and may become meaningless. In the server device 100 according to the present embodiment, even if the default number of replicas is inappropriate, the optimum number of replicas can be finally applied.
(System configuration) Next, the system configuration of the server device that realizes the above-mentioned data processing will be described. As shown in FIG. 1, access to the data distribution system 200 is performed via the load balancer 110. The load balancer 110 collectively receives requests 120 input from the user using the client terminal, and transmits the received data and processing requests to any server device 100 constituting the data distribution system 200.
The allocation operation of data and processing requests by the load balancer 110 is not particularly limited. For example, the received ones may be sequentially allocated in the order of the device number of the server device 100, or the operating status of each server device 100 may be monitored and busy. It may be randomly assigned to other than the server device 100 in the state.
Next, the configuration of each server device constituting the data distribution system 200 will be described. Each server device 100 has a configuration including an application execution unit 101, a cooperation processing unit 102, and a storage unit 103. Then, when accessing each data, the application executed by the application execution unit 101 is set to operate on all the server devices 100 in which the corresponding data is arranged. Further, the storage unit 103 is a recording area in which data is actually arranged, and is realized by various memories and disks. Since a known technique is used for the actual data writing process to the storage unit 103, the description thereof is omitted here.
The feature of the server device 100 according to the present embodiment is the cooperation processing unit 102. When the request 120 is newly allocated by the load balancer 110, the application execution unit 101 converts the input request 120 into an access process (request) for the actual data and sends it to the cooperation processing unit 102 of the own device. Therefore, when the cooperation processing unit 102 receives the request 120, the cooperation processing unit 102 selects a server device (hereinafter, referred to as an arrangement server) for arranging the data according to the number of replicas set for the target data included in the request 120. Then, the request 120 is sent to the cooperation processing unit 102 of the selected server device 100.
In addition, the linkage processing unit 102 provides information for identifying data (for example, when the number of replicas is set for newly input data or when the number of replicas of certain data is updated by a function described later). The number of replicas corresponding to the data is transmitted to each server device 100 together with the file name including the data contents). In this way, by holding the information on the number of replicas set for each data by the cooperation processing unit 102, the server device 100 identifies the placement server of the target data regardless of which data the request 120 is allocated. , Appropriate processing can be executed.
(Hardware configuration of server device) Next, a specific hardware configuration of the server device will be described. FIG. 2 is a block diagram showing a hardware configuration of the server device according to the present embodiment. In FIG. 2, the server device 100 includes a CPU (Central Processing Unit) 201, a ROM (Read-Only Memory) 202, a RAM (Random Access Memory) 203, a magnetic disk drive 204, a magnetic disk 205, and communication I. It includes / F (Interface) 206, an input device 207, and an output device 208. Further, each component is connected by a bus 210.
Here, the CPU 201 controls the entire server device 100. The ROM 202 stores various programs such as a boot program and a data processing program for realizing the data processing according to the present embodiment. RAM203 is used as a work area for CPU201. The magnetic disk drive 204 controls the update / reference of data to the magnetic disk 205 according to the control of the CPU 201. The magnetic disk 205 stores data written under the control of the magnetic disk drive 204. In the hardware configuration of FIG. 2, the magnetic disk 205 is used as the recording medium that plays the role of the storage unit 103, but other recording media such as an optical disk and a flash memory may be used.
The communication I / F 206 is connected to a network (NET) 209 such as LAN (Local Area Network), WAN (Wide Area Network), and the Internet through a communication line, and is connected to another server device 100 or load balancer 110 via this network 209. Be connected. The communication I / F 206 controls the internal interface with the network 209 and controls the input / output of data from the external device. As a configuration example of the communication I / F 206, for example, a modem or a LAN adapter can be adopted.
The input device 207 receives an external input to the server device 100. Specific examples of the input device 207 include a keyboard and a mouse. As illustrated in FIG. 1, since the server device 100 executes the application in response to the request 120 allocated by the load balancer 110, the request 120 is not input from the input device 207, and the server device 100 is not input. It is prepared for the purpose of maintenance and management of.
In the case of a keyboard, for example, it is provided with keys for inputting characters, numbers, various instructions, etc., and data is input. Further, it may be a touch panel type input pad, a numeric keypad, or the like. In the case of a mouse, for example, move the cursor, select a range, move the window, and resize the window. Further, a trackball, a joystick, or the like may be used as long as it has the same function as a pointing device.
The output device 208 outputs the data arranged in the server device 100, the execution status of the application, the access pattern of each arranged data, and the analysis result thereof. Specific examples of the output device 208 include a display and a printer.
In the case of a display, for example, data such as a cursor, an icon, a toolbox, a document, an image, and functional information is displayed. Further, as this display, a CRT, a TFT liquid crystal display, a plasma display, or the like can be adopted. In the case of a printer, for example, image data and document data are printed. Further, a laser printer or an inkjet printer can be adopted.
The above-mentioned input device 207 and output device 208 are not essential configurations due to the characteristics of the server device 100, and the configurations may be appropriately changed according to the convenience of the administrator.
(Functional configuration of the cooperative processing unit) Next, the detailed processing of the cooperative processing unit 102 will be described. As described with reference to FIG. 1, bidirectional communication is possible between the cooperation processing units 102 of each server device 100, and when a request 120 for a certain data is assigned, it is controlled by the cooperation processing unit 102 described later. Request 120 (information that identifies the type of processing and data and data content) is sent and received. In addition, when the number of replicas for a certain data is newly determined or changed, new replica number information (information for identifying data, data contents, new number of replicas) is transmitted and received.
Therefore, the functional configuration for the cooperative processing unit 102 of the server device 100 according to the present embodiment to realize the above-mentioned control will be described below. FIG. 3 is a block diagram showing the functional configuration of the cooperative processing unit. As shown in FIG. 3, the cooperation processing unit 102 includes an acquisition unit 301, a selection unit 302, a transmission / reception unit 303, a setting unit 304, a judgment unit 305, an execution unit 306, and a determination unit 307. .. Specifically, the function serving as the control unit (acquisition unit 301 to determination unit 307) causes the CPU 201 to execute a program stored in a storage area such as ROM 202, RAM 203, or magnetic disk 205 shown in FIG. The function is realized by this or by communication I / F206.
When the request 120 is input as a processing request for arbitrary data, the acquisition unit 301 acquires the number of replicas set in this data. As described above, if the data is already arranged in the 100 million copies 103 of any of the server devices 100 constituting the data distribution system 200, all the server devices 100 have information on the number of replicas of this data. It is being held. Therefore, the acquisition unit 301 acquires the number of replicas set in the target data of the input request 120. The correspondence information between the number of replicas acquired by the acquisition unit 301 and the data is stored in a storage area such as 100 million copies 103 (for example, magnetic disk 305).
The selection unit 302 selects the server device 100 in which the data targeted by the request 120 is arranged from the server devices 100 constituting the data distribution system 200 for the number of replicas by using a predetermined algorithm. The algorithm used as a selection criterion here is not particularly limited. For example, the hash value of the input data is obtained, and the server device 100 for the number of replicas is selected as the placement server, starting with the server device 100 whose device number matches the remainder when divided by the number of servers.
In addition, if the number of server devices 100 is N, the identification number of the input data is converted to an N-ary number, and the server device 100 that matches the converted identification number and the device number is the first. A method of selecting as many server devices as the number of replicas as the placement server may be used. The information of the placement server selected by the selection unit 302 in this way is stored in a storage area such as the storage unit 103.
The transmission / reception unit 303 performs two-way communication with the cooperation processing unit 102 of the other server device 100. For example, when the transmission / reception unit 303 acquires an arbitrary number of replicas of data by the acquisition unit 301, the transmission / reception unit 303 transmits the data to all the server devices 100 constituting the data distribution system 200. Further, the transmission / reception unit 303 transmits the request 120 to each server device 100 for the number of replicas selected as the placement server by the selection unit 302. The transmission / reception unit 303 also plays a role of receiving the number of replicas transmitted from the cooperation processing unit 102 of the other server device 100 and the request 120.
The setting unit 304 sets a value previously given as the number of replicas of this data when a request for new write processing is input as the request 120 for arbitrary data. In the acquisition unit 301, the number of replicas can be acquired from the storage unit 103 if the data is already arranged in any of the server devices 100, but in the case of new writing, the number of replicas is not held. Therefore, the initial value set in advance by the administrator of the data distribution system 200 can be set as the number of replicas. At this time, even if the number of replicas unsuitable for the data access content is set, the number of replicas is dynamically changed by the determination unit 307 described later, so that the processing efficiency is not adversely affected.
Even if the number of replicas of arbitrary data is set by the setting unit 304, the selection unit 302 and the transmission / reception unit 303 select the placement server and perform replicas in the same manner as when the number of replicas is set by the acquisition unit 301. Send numbers and requests 120.
When the new write request 120 is input and the number of replicas of arbitrary data is set by the setting unit 304, the determination unit 305 determines whether or not the number of replicas is equal to the total number of the server device 100 groups. When the determination unit 305 determines that the number of replicas and the total number of the server device 100 groups are equal, the transmission / reception unit 303 sends a request 120 indicating a request for new write processing of arbitrary data to all the server device 100 groups. To do.
Then, when the determination unit 305 determines that the number of replicas and the total number of the server device 100 groups are not equal, the selection unit 302 should arrange this data from the server device 100 groups using the above-mentioned algorithm. Select the server device 100 for the number of replicas. Then, the transmission / reception unit 303 transmits a request 120 indicating an update process for this data to each server device 100 for the number of replicas selected by the selection unit 302.
Further, when a request for write processing for this data arranged in any of the server device 100 groups is input as a request 120 for a certain data, the judgment unit 305 sets the number of replicas of this data and the server device 100 group. Determine if the total number of is equal to or not. When the determination unit 305 determines that the number of replicas and the total number of the server device 100 groups are equal, the transmission / reception unit 303 transmits a request 120 indicating a request for the above data update process to all the server device 100 groups.
Then, when the determination unit 305 determines that the number of replicas and the total number of the server device 100 groups are not equal, the selection unit 302 arranges this data from the server device 100 groups using the above-mentioned algorithm. Select the server device 100 for the number of replicas. Even when such processing is performed, the transmission / reception unit 303 transmits a request 120 indicating a writing process for this data to each server device 100 for the number of selected replicas.
When a request for reference processing is input as a request 120 for a certain data, the transmission / reception unit 303 refers to any one of the server devices 100 for the number of replicas selected by the selection unit 302. Request 120 representing the request of is sent.
Further, when the execution unit 306 receives the request 120 transmitted from the own device or another server device via the transmission / reception unit 303, the execution unit 306 executes the process according to the request 120. Therefore, if the update request is accepted, the update process is performed, and if the reference request is accepted, the reference process is performed. When the execution unit 306 receives the update request as the request 120, the execution unit 306 goes into a standby state when the update process is completed as it is. On the other hand, when the execution unit 306 receives the reference request as the request 120, the execution unit 306 returns the result of the reference processing to the target data to the server device 100 that is the source of the request 120 via the transmission / reception unit 303.
This is because, from the viewpoint of the server device 100 to which the request 120 is assigned, the execution of the request 120 is processed as if it were performed by the own device. Therefore, if the reference result is not returned to the server device 100 that is the source of the processing request, it is determined that the processing request has not been executed in the server device 100 to which the request 120 is assigned, so that the reference result is returned. By doing so, such a situation can be prevented.
Further, the server device 100 may have a function of independently determining the number of replicas inside the device. The determination unit 307 can determine the number of replicas of each data by referring to the request 120 for each data executed by the execution unit 306 at any timing.
Therefore, the transmission / reception unit 303 transmits a new number of replicas to all the server devices 100 even when the determination unit 307 determines the number of replicas different from the currently set number of replicas. The selection unit 302 of each server device 100 selects the placement server based on the number of new replicas transmitted. The execution unit 306 writes the target data to the storage unit 103 when the own device is newly selected by the selection unit 302 as the placement server. On the other hand, when the selection unit 302 does not newly select the target data as the placement server, that is, when the selection of the placement server is omitted, the target data is deleted from the storage unit 103.
Further, in the determination unit 307, when the average value of the processing time for the target data executed by the execution unit 306 is equal to or more than a predetermined value, the total number of the server device 100 groups is determined as the number of replicas of the target data, and the target data is processed. When the average value of time is less than a predetermined value, the minimum value that can be set as the number of replicas can be determined as the number of replicas.
To specifically explain the comparison between the average processing time for the target data and the predetermined value described above, for example, the average processing time for the target data when the number of replicas of the target data is the total number of 100 groups of server devices. A method of determining the number of replicas by obtaining the average value of the value and the processing time for the target data when the number of replicas of the target data is 1 and comparing the two can be mentioned. By this comparison, if the former is smaller than the latter, the total number of 100 server devices can be determined as the number of replicas, and if the latter is smaller than the former, the minimum value that can be set as the number of replicas can be determined as the number of replicas.
As described above, in the server device 100 according to the present embodiment, the request 120 received by the load balancer 110 is set as the target data of the request 120 regardless of which server device 100 in the data distribution system 200 is assigned. The placement destination of the target data can be determined according to the number of replicas made. Therefore, no matter what kind of content the request 120 is received, it is possible to efficiently distribute the data by selecting the placement server according to the number of replicas. Hereinafter, the specific operations of the server device 100 having the above-mentioned functions will be sequentially described separately at the time of data update and at the time of data reference.
(Data update process) First, the data update process in the server device 100 will be described. FIG. 4 is an explanatory diagram showing data update processing in the server device. A procedure for the cooperation processing unit 102 of each server device 100 to access the target data according to the number of replicas set for each data will be described with reference to FIG.
First, the load balancer 110 allocates the process of updating the target data (writing to new data, writing to the placed data, deleting the placed data) in one of the server devices 100 (server device 100-1 in FIG. 4). Is done. The cooperative processing unit 102 determines the server for arranging the data by a predetermined algorithm.
Here, the procedure of the determination process of the placement server will be described. FIG. 5 is a flowchart showing the procedure of the placement server selection process. In the flowchart of FIG. 5, first, the number of newly determined or updated replicas is acquired (step S501). Then, it is determined whether or not the acquired number of replicas is equal to the number (total number) of the server devices 100 (step S502). Here, the number of server devices 100 to be compared with the number of replicas means the total number of server devices 100 set so that request 120 is allocated by the load balancer 110. Therefore, in the examples shown in FIGS. 1 and 4, the number of server devices 100 is n.
If it is determined in step S502 that the number of replicas and the number of servers are equal (step S502: Yes), an update request must be made to all server devices 100, so all server devices 100 are determined to be the placement servers. (Step S503), a series of placement server selection processes is terminated.
On the other hand, if it is determined in step S502 that the number of replicas and the number of server devices 100 are not equal (step S502: No), the server devices 100 for the number of replicas are determined as the placement server from all the server devices 100. Move to the process for. First, the hash value of the target data is calculated (step S504). Then, it is determined whether or not the number of replicas is 1 (step S505).
If the number of replicas is determined to be 1 in step S505 (step S505: Yes), the remainder when the calculated hash value is divided by the number of servers (total number) is calculated, and the remainder value matches the device number. The server device 100 to be used is selected as the placement server (step S506), and a series of processes is completed. On the other hand, when it is determined that the number of replicas is other than 1 (step S505: No), the remainder when the hash value is divided by the number of servers is obtained as in step S506, and the device number matching this remainder value is given. Starting with the server device 100, the server device 100 for the number of replicas is selected as the placement server (step S507), and the series of processes is completed.
To specifically explain the processes in steps S506 and S507, for example, the remainder (odd) obtained by dividing the hash value by the number of servers is obtained. Then, if the number of replicas = 1, as in step S506, the server device 100 whose odd is set to the device number is selected as the placement server. Then, when the number of replicas = m (1 <m <n), m server devices in which odd, odd + 1, ..., odd + m-1 are set as device numbers as in step S507. 100 is the deployment server.
Next, returning to FIG. 4, the processing of the update processing patterns (1) and (2) will be described. When the cooperation processing unit 102 determines the placement server, the received update request must be transmitted to the placement server. At this time, in the cooperation processing unit 102, when the own device is the placement server, the update processing pattern (1) that writes in response to the update request and the update that writes in response to the update request to the placement server other than the own device. Processing pattern (2) is required.
FIG. 6-1 is a flowchart showing a procedure of data update processing (in the case of a server device that receives an update request from the application execution unit of the own device). In the flowchart of FIG. 6-1 first, it is determined whether or not an update request has been received from the application execution unit 101 of the own device (step S611). Here, the state waits until the data update request is received (step S611: No loop), and when the data update request is received (step S611: Yes), the placement server selection process is performed (step S612).
Next, it is determined whether or not the own device is the placement server by referring to the result of the placement server selection process in step S612 (step S613). Here, if it is determined that the local device is the deployment server (step S613: Yes), the target data is updated in response to the data update request (step S614), and the update request is sent to another deployment server (step S614). S615), the series of processing is completed. If it is determined in step S613 that the own device is not the deployment server (step S613: No), the process proceeds to step S615 and an update request is sent to another deployment server without making the update request in step S614. Send (step S615) to end the series of processes.
Then, when the update request is transmitted to the other server device 100 determined by the placement server as in step S615, the destination server device 100 (in the case of FIG. 4, the server device 100-n) further Processing is required. Figure 6-2 is a flowchart showing the procedure of data update processing (in the case of a server device in which an update request is received from the cooperation processing unit of another server device). In the flowchart of FIG. 6-2, first, it is determined whether or not an update request has been received from another server device 100 (step S621). Here, it waits until the update request is received (step S621: No loop), and when it is determined that the update request has been received (step S621: Yes), the target data is updated according to the update request (step S622). Ends a series of processes.
(Data reference processing) Next, the data reference process in the server device 100 will be described. FIG. 7 is an explanatory diagram showing data reference processing in the server device. Also in the case of FIG. 7, it is assumed that request 120 (reference request) is assigned to the server device 100-1. The cooperative processing unit 102 determines the allocation server of the data by a predetermined algorithm as in the case of the update processing (see FIG. 5).
When the cooperation processing unit 102 determines the placement server, the cooperation processing unit 102 transmits the received reference request to the placement server. At this time, when the own device is the placement server, the cooperation processing unit 102 transmits the processing pattern (1) for reading according to the reference request and the reference request to the placement server other than the own device, and sends the request. A processing pattern (2) in which the received server device 100 reads according to the request is required. Further, in the case of the above processing pattern (2), the server device 100-n to which the reference request has been transmitted must transmit this reference result to the server device 100-1 of the transmission source. Therefore, the processing in each of the processing patterns (1) and (2) will be described below.
FIG. 8-1 is a flowchart showing the procedure of data reference processing (in the case of a server device that receives a reference request from the application execution unit of the own device). In the flowchart of FIG. 8-1, first, it is determined whether or not the request 120 (reference request) is received from the application execution unit 101 of the own device (step S811). Here, the state waits until the data reference request is received (step S811: No loop), and when the data reference request is received (step S811: Yes), the deployment server selection process is performed (step S812).
Next, it is determined whether or not the number of replicas of the target data is 1 (step S813). When the number of replicas is set to 1 (step S813: Yes), the reference server for referencing the server device 100 determined as the placement server by the placement server selection process in step S812 in response to the reference request. (Step S814). That is, it means that there is no server device 100 in which the target data is arranged other than the determined arrangement server.
On the other hand, when the number of replicas is set to a value other than 1 (step S813: No), any one of the server devices 100 determined as the placement server by the placement server selection process in step S812 responds to the reference request. Determine the reference server for data reference (step S815).
After the reference server is determined, the own device then determines whether or not it is the deployment server determined in step S812 (step S816). Here, when it is determined that the local device is the placement server (step S816: Yes), the target data stored in the storage unit 103 of the local device is referred to (step S817), and a series of processing is terminated. .. On the other hand, if it is determined that the local device is not the placement server (step S816: No), a reference request is sent to the reference server (step S818), and the series of processes is terminated.
Then, when the reference process is transmitted to the other server device 100 determined as the reference server as in step S818, the destination server device 100 (server device 100-n in the case of FIG. 7) further performs further. Processing is required. Figure 8-2 is a flowchart showing the procedure of data reference processing (in the case of a server device that receives a reference request from the cooperative processing unit of another server device).
In the flowchart of FIG. 8-2, first, it is determined whether or not a reference request has been received from the cooperation processing unit 102 of the other server device 100 (step S821). Here, it waits until the reference request is received (step S821: No loop), and when it is determined that the reference request has been received (step S821: Yes), the target data is referenced according to the reference request (step S822). The reference result is transmitted to the server device 100 that is the source of the reference request (step S823), and the series of processes is completed.
(Data relocation processing when changing the number of replicas) Next, the data rearrangement process when the number of replicas is changed will be described. In the server device 100 according to the present embodiment, it is necessary to newly arrange the target data in the storage unit 103 or delete the arranged target data according to the change in the number of replicas. Hereinafter, the processing contents will be described separately for each case.
FIG. 9 is a flowchart showing a data rearrangement process when the number of replicas is changed. In the flowchart of FIG. 9, first, it is determined whether or not a predetermined time has elapsed (step S901). Here, the process waits until the predetermined time elapses (step S901: No loop), and when the predetermined time elapses (step S901: Yes), the number of replicas is determined (step S902).
Note that this step S901 is a process of determining the update timing of the number of replicas on a time basis. As already mentioned, the update timing setting is arbitrary and can be set freely. Therefore, for example, the update timing may be set based on the number of processes, such as determining for specific data or whether or not a predetermined number of request processes are allocated to the own device.
Further, in the replica number determination process in step S902, various access pattern analysis tools may be used, or settings from the administrator may be accepted. Further, the number of replicas may be determined independently by a procedure as described later.
Once the number of replicas has been determined, it is then determined whether or not the number of replicas has been changed by the determination in step S902 (step S903). Here, if the number of replicas has not changed (step S903: No), the number of deployed servers does not change, so the series of processes ends as it is.
On the other hand, if the number of replicas is changed (step S903: Yes), the information on the changed number of replicas is transmitted to the other server device 100 (step S904). Then, it is determined whether or not the number of replicas has increased due to the change (step S905). If the number of replicas has increased (step S905: Yes), the data relocation process at the time of increase (step S906) has been performed, if the number of replicas has not increased (step S905: No), that is, if the number of replicas has decreased. Is to perform a data rearrangement process at the time of decrease (step S907), and end a series of processes.
When the number of replicas increases First, the process when the number of replicas increases will be described. FIG. 10 is a flowchart showing a data rearrangement process when the number of replicas increases. In the flowchart of FIG. 10, first, the placement server selection process using the changed number of new replicas is performed (step S1001). The placement server selection process performed in step S1001 is the placement server selection process described in FIG. The deployment server selected in step S1001 is used as the new deployment server.
Next, it is determined whether or not the number of old replicas set for the target data = 1 (step S1002). Here, if it is determined that the number of old replicas = 1 (step S1002: Yes), the old placement server determined based on the number of old replicas is definitely determined to be the new placement server, so this placement Determine the server as the data source (step S1003). This data transmission source means the server device 100 in which the master data is arranged, which transmits the target data to the server device 100 newly added to the arrangement server.
On the other hand, if it is determined that the number of old replicas is not = 1 (step S1002: No), one of the old placement servers, which is the placement server determined by the number of old replicas, is determined as the data transmission source (step S1004). ). Again, if the number of replicas increases, the deployment server determined by the number of old replicas is definitely included in the deployment server determined by the number of new replicas, so one of the old deployment servers is determined as the data source. do it.
After that, it is determined whether or not the own device is the data transmission source (step S1005), and if the own device is determined as the data transmission source (step S1005: Yes), the target data is sent to the server device 100 that has become the newly deployed server. Is sent (step S1006) to end the series of processes. On the other hand, if the own device is not the data transmission source (step S1005: No), the series of processes is terminated as it is. In this case, since the server device 100 determined as the data transmission source arranges the data, the own device does not have to do anything.
When the number of replicas decreases Next, the process when the number of replicas decreases will be described. FIG. 11 is a flowchart showing a data rearrangement process when the number of replicas decreases. In the flowchart of FIG. 11, first, the placement server selection process using the number of old replicas is performed (step S1101). The placement server selection process performed in step S1101 is the placement server selection process described with reference to FIG.
With reference to the placement server selection process in step S1101, it is determined whether or not the local device is the placement server in the number of old replicas (step S1102). Here, if it is determined that the own device is not the placement server (step S1102: No), the target data whose number of replicas has changed this time is not placed in the own device, so the series of processing is terminated as it is. To do.
On the other hand, in step S1102, when it is determined that the own device is the placement server (step S1102: Yes), this time, the placement server selection process using the number of new replicas is performed (step S1103). Then, referring to the placement server selection process in step S1103, it is determined whether or not the own device is the placement server in the number of new replicas (step S1104).
If it is determined in step S1104 that the own device is the placement server in the number of new replicas (step S1104: Yes), the target data placed in the own device is retained, so the series of processes is terminated as it is. .. On the other hand, when it is determined that the own device is not the placement server in the new number of replicas (step S1104: No), the target data placed in the own device is deleted (step S1105), and the series of processes is terminated as it is.
As described above, the server device 100 according to the present embodiment can efficiently arrange the target data in response to the dynamic change in the number of replicas.
(Processing to determine the number of replicas) Next, the process of determining the number of replicas will be described. As described above, the method for setting the number of replicas in the server device 100 according to the present embodiment is not uniform. For example, the administrator of the data distribution system 200 may analyze the access pattern and set the number of replicas of each data based on the analysis result, or prepare a tool for analyzing the access pattern and analyze the result by this tool. The number of replicas may be set by.
However, by providing the server device 100 with a function of independently determining the number of replicas, the burden on the administrator can be reduced. Therefore, here, a specific example will be described when the number of replicas is automatically determined by the cooperation processing unit 102 of each server device 100.
The cooperation processing unit 102 of the server device 100 calculates the performance of data access processing for each data in consideration of the reference / update ratio of a certain data (breakdown of access to the data). The information used in this calculation process is the number of replicas, the writing time to the storage unit 103 or the reading time from the storage unit 103, and the communication time between the server devices 100. In an environment where the ratio of reference / update requests changes, it is necessary to periodically calculate the performance using the ratio of reference / update requests at that time to determine the number of replicas with the highest performance. The procedure for determining the number of replicas using the above information will be described below.
Before the explanation, the variables used are listed below.
Update rate: W [%] Reference ratio: R [%] Number of servers: N [units] Number of replicas: r (integer of r> 0) Communication time between server devices: Tt [sec] Write time: Tw [sec] Read time: Tr [sec] Average latency at update: Lw Average latency at reference: Lr
First, the cooperation processing unit 102 counts reference / update requests for each data at a certain time interval, and obtains a reference ratio R / update ratio W for each data. Specifically, if 180 reference processes are performed and 20 update processes are performed in one hour, the reference ratio R = 90 [%] and the update ratio W = 10 [%]. .. Here, as an example, the reference ratio R and the update ratio W are obtained from the number of reference / update requests to the data generated at a certain time interval, but what is the reference and update of each of the certain number of requests? The reference ratio R and the update ratio W may be obtained based on the number of times. For example, if 90 out of 100 requests are references and 10 are updates, then R = 90 [%] and W = 10 [%].
Next, the cooperation processing unit 102 calculates the above-mentioned reference ratio R and update ratio W, and at the same time, calculates the average latency of each of the data reference and update according to the number of replicas r, and the average latency becomes the lowest. Determine the number of replicas. The procedure for determining the number of replicas will be described in detail below with reference to FIGS. 12 and 13.
The communication time Tt, write time Tw, and read time Tr between the server devices 100 may be obtained by averaging several measured values. Alternatively, the communication time Tt, write time Tw, and read time Tr between the actual server devices 100 may be measured individually, given as specifications in advance, or the test values are disclosed. If it is known in advance, the value may be used. Here, the means for obtaining these values is not particularly limited, and the administrator of the data distribution system 200 can appropriately select them. Then, when an update request is transmitted from one server device 100 to a plurality of other server devices 100, the transmission and the actual update request are sequentially performed.
Average latency when updating data First, the average latency at the time of data update will be described. FIG. 12 is an explanatory diagram showing a latency calculation procedure at the time of data update. As shown in FIG. 12, the average latency at the time of data update needs to be calculated in consideration of the latency in each case when the own device is determined to be the placement server and the case where the own device is not the placement server.
First, when the server device 100 that received the update request is the placement server for the target data (own device = placement server), the destination of the data update request is a total of r-1 remote server devices other than the own device. , The actual update process is performed by n server devices 100 including the own device. Therefore, the latency Lwa in such a case can be obtained by the following equation (1).
Lwa = r * Tw + (r-1) * Tt ... (1)
On the other hand, if the server device 100 that received the update request is not the placement server for the target data (own device placement server), the destination of the data update request is a total of r units of remote server devices 100 other than the own device. The actual update process is also performed on r servers. Therefore, the latency Lwb in such a case can be obtained by the following equation (2).
Lwb = r * (Tw + Tt) ... (2)
The probability that the server device 100 that receives the update request will be the placement server for the target data is r / N, and the probability that it will not be the placement server is (Nr) / N. Therefore, the average latency Lw at the time of updating is as follows (3). ) Is calculated by the formula.
Lw = Lwa * r / N + Lwb * (Nr) / N = {r * Tw + (r-1) * Tt} * r / N + r * (Tt + Tw) * (Nr) / N ... (3)
Average latency when referencing data Next, the average latency at the time of data reference will be described. FIG. 13 is an explanatory diagram showing a latency calculation procedure at the time of data reference. As shown in FIG. 13, the average latency at the time of data reference also needs to be calculated in consideration of the latency in each case when the own device is determined to be the placement server and the case where the own device is not the placement server.
First, when the server device that received the reference request is the placement server for the target data (own device = placement server), it is not necessary to send the received data reference processing to another server device 100, so the latency is the reference in the own device. Only the processing time. Therefore, the latency Lra in such a case can be obtained by the following equation (4).
Lra = Tr ... (4)
On the other hand, if the server device that received the reference request is not the placement server for the target data (own device placement server), the data reference process is sent to one server device 100 in the placement server, and that one server device. Reference processing is performed at 100. Therefore, the latency Lrb in such a case can be obtained by the following equation (5).
Lrb = Tr + Tt ... (5)
Also in the case of reference processing, the probability that the server device 100 that received the reference request will be the placement server for the target data is r / N, and the probability that it will not be the placement server is (Nr) / N, so the average latency Lr at the time of reference is as follows. It is calculated by equation (6).
Lr = Lwa * r / N + Lwb * (Nr) / N = Tr * r / N + (Tt + Tr) * (Nr) / N ... (6)
As explained above, since the data reference / update ratios R and W are obtained by monitoring for a certain period of time as described above, the average latency L when the number of replicas r is used is as follows (using that ratio). It is calculated by equation 7).
L = Lw * W / 100 + Lr * R / 100 = {{r * Tw + (r-1) * Tt} * r / N + r * (Tt + Tw) * (Nr) / N} * W / 100 + {Tr * r / N + (Tt + Tr) * (Nr) / N} * R / 100 ... (7)
Since the above formula for obtaining the average latency L is a linear expression of the number of replicas r, the average latency L is minimized when the number of replicas r = 1 or when the number of replicas r = the number of servers N. Is it? Therefore, only the average latency L1 when r = 1 and the average latency LN when r = N are calculated, and the number of replicas whose average latency is lower is used as the new number of replicas. That is, the number of replicas is uniquely determined by not calculating the average latency for all r from 1 to N, but calculating only the average latency when r = 1 and r = N as shown below. can do.
When L1 <LN, set the number of replicas to 1. When L1> LN, set the number of replicas to N
If the number of replicas is set to 1 during the actual operation of the data distribution system 200, it may not be preferable from the viewpoint of availability. Therefore, when the decision to "set the number of replicas to 1" is made, a method of setting the minimum number of replicas set in advance by the administrator of the data distribution system 200 may be used.
As described above, according to the present embodiment, by monitoring not only the data reference frequency but also the update frequency, it is possible to detect a case where the processing efficiency decreases as the number of replicas increases. In such a case, it is possible to avoid reducing the efficiency of the update process by deleting unnecessary replicas. In addition, the optimum number of replicas can be set individually for each data without being affected by the number of replicas of other data, and the data can be arranged efficiently.
The data processing method described in the present embodiment can be realized by executing a program prepared in advance on a computer such as a personal computer or a workstation. This program is executed by being recorded on a computer-readable recording medium such as a hard disk, flexible disk, CD-ROM, MO, or DVD, and read from the recording medium by the computer. In addition, this program may be a medium that can be distributed via a network such as the Internet.
Further, the server device 100 described in the present embodiment is a PLD (Programmable Logic) such as an application specific integrated circuit (IC (hereinafter, simply referred to as ASIC) such as a standard cell or a structured ASIC (Application Specific Integrated Circuit), or an FPGA. It can also be realized by Device). Specifically, for example, the functions (acquisition unit 301 to determination unit 307) of the cooperation processing unit 102 of the server device 100 described above are defined by the HDL description, and the HDL description is logically synthesized and given to the ASIC or PLD. Therefore, the server device 100 can be manufactured.
The following additional notes are further disclosed with respect to the above-described embodiment.
(Appendix 1) Computers that make up a group of computers that can communicate with each other, When a processing request for arbitrary data is input, an acquisition means for acquiring the number of duplicates set for the arbitrary data, A selection means for selecting a computer as a destination for arranging arbitrary data from the computer group for the number of duplicates using a predetermined algorithm. A copy number transmission means for transmitting the number of copies of the arbitrary data acquired by the acquisition means to all the computers. A processing request transmitting means for transmitting the processing request to each computer for the number of duplicates selected by the selection means. A data processing program characterized by functioning as.
(Appendix 2) The computer is further added. When a request for processing new writing of the arbitrary data is input as a processing request for the arbitrary data, the function is made to function as a setting means for setting a value given in advance as the number of duplicates of the arbitrary data. , When the number of duplicates of the arbitrary data is set by the setting means, the selection means selects the computer to which the arbitrary data is arranged from the computer group by using a predetermined algorithm. Select minutes and The copy number transmitting means transmits the copy number of the arbitrary data set by the setting means to the computer group, and transmits the copy number transmission means. The data processing program according to Appendix 1, wherein the processing request transmitting means transmits a request for new writing processing of the arbitrary data to each computer for the number of duplicates selected by the selecting means. ..
(Appendix 3) Further, the computer When the number of duplicates of the arbitrary data is set by the setting means, it functions as a determination means for determining whether or not the number of duplicates is equal to the total number of the computer group. When the determination means determines that the number of duplicates and the total number of the computer groups are equal, the processing request transmitting means transmits a request for new writing processing of the arbitrary data to all the computer groups. The data processing program described in Appendix 2, which is characterized by the above.
(Appendix 4) When the request for write processing for the arbitrary data arranged in any of the computer groups is input as the processing request for the arbitrary data, the determination means of the arbitrary data. Determine if the number of duplicates is equal to the total number of computers. When the determination means determines that the number of duplicates and the total number of the computer group are not equal, the selection means arranges the arbitrary data from the computer group by using the predetermined algorithm. Select the computers that are used for the number of duplicates, The data processing program according to Appendix 3, wherein the processing request transmitting means transmits a request for writing processing to the arbitrary data to each computer for the number of duplicates selected by the selecting means.
(Appendix 5) The processing request transmitting means is the selection means when a request for reference processing for the arbitrary data arranged in any of the computer groups is input as a processing request for the arbitrary data. The data processing program according to any one of Supplementary note 1 to 4, wherein a request for reference processing for the arbitrary data is transmitted to any one of the computers selected by the number of duplicates. ..
(Appendix 6) Further, the computer Described in any one of Appendix 1 to 5, which is characterized in that when a processing request transmitted from its own device or another computer is received, it functions as an execution means for executing the processing according to the processing request. Data processing program.
(Appendix 7) When the execution means receives the transmission of the request for reference processing for the arbitrary data as the processing request, the result of the reference processing for the arbitrary data executed in response to the processing request is displayed. The data processing program according to Appendix 6, wherein the processing request is returned to the computer that is the source of the processing request.
(Appendix 8) Further, the computer It functions as a determination means for determining the number of duplicates of the arbitrary data by referring to the processing request for the arbitrary data executed by the execution means at any arbitrary timing. When the number of copies determined by the determination means is different from the number of copies currently set, the means for transmitting the number of copies transmits the determined number of copies to all the computer groups. When the determined number of copies is transmitted by the copy number transmission means, the selection means newly sets a computer for arranging the arbitrary data from the group of computers according to a predetermined algorithm. Select the number of duplicates determined above, The execution means is a computer that writes the arbitrary data when the own device is newly selected by the selection means to the computer on which the arbitrary data is newly arranged, and newly arranges the data by the selection means. The data processing program according to Appendix 6 or 7, wherein the data is deleted when the data is no longer selected.
(Appendix 9) When the average value of the processing time for the arbitrary data is a predetermined value or more when the total number of the computer groups is the number of duplicates of the arbitrary data, the determination means determines the total number of the computer groups. The data processing program according to Appendix 8, wherein the number of duplicates of arbitrary data is determined.
(Appendix 10) When the average value of the processing time for the arbitrary data is less than a predetermined value when the number of duplicates of the arbitrary data is 1, the determination means determines the number of duplicates of the arbitrary data. The data processing program according to Appendix 8 or 9, wherein the minimum value that can be set is determined as the number of duplicates.
(Appendix 11) A server device that constitutes a group of server devices that can communicate with each other. When a processing request for arbitrary data is input, an acquisition means for acquiring the number of duplicates set for the arbitrary data, and an acquisition means. A selection means for selecting the server device to which the arbitrary data is to be arranged from the server device group for the number of duplicates by using a predetermined algorithm. A copy number transmission means for transmitting the copy number of the arbitrary data acquired by the acquisition means to all the server device groups, and a copy number transmission means. A processing request transmitting means for transmitting the processing request to each server device for the number of duplicates selected by the selecting means, and When a processing request sent from the own device or another server device is received, an execution means for executing the processing according to the processing request and an execution means. A determination means for determining the number of duplicates of the arbitrary data by referring to a processing request for the arbitrary data executed by the execution means at any timing is provided. When the number of copies determined by the determination means is different from the number of copies currently set, the means for transmitting the number of copies transmits the determined number of copies to all the server devices. When the determined number of copies is transmitted by the copy number transmission means, the selection means newly uses a server device for arranging the arbitrary data from the server device group as a predetermined algorithm. Depending on the number of duplicates determined above, The executing means writes the arbitrary data when the own device is newly selected by the selection means to the server device to newly arrange the arbitrary data, and newly arranges the data by the selection means. A server device characterized in that the data is deleted when the server device is no longer selected.
(Appendix 12) Computers that make up a group of computers that can communicate with each other When a processing request for arbitrary data is input, an acquisition process for acquiring the number of duplicates set for the arbitrary data, and an acquisition process. A selection step of selecting the computer to which the arbitrary data is to be placed from the computer group for the number of duplicates by using a predetermined algorithm. A copy number transmission step of transmitting the copy number of the arbitrary data acquired by the acquisition step to all the computer groups, and a copy number transmission step. A process request transmission step of transmitting the process request to each computer for the number of duplicates selected by the selection step, and a process request transmission step. When a processing request sent from the own device or another computer is received, an execution process that executes the processing according to the processing request and an execution process. At any timing, a determination step of determining the number of duplicates of the arbitrary data by referring to the processing request for the arbitrary data executed by the execution step is executed. further, In the copy number transmission step, when a copy number different from the copy number currently set by the determination step is determined, the determined copy number is transmitted to all the computer groups. In the selection step, when the number of copies determined by the number of copies transmission step is transmitted, a computer for arranging the arbitrary data from the group of computers is newly selected according to a predetermined algorithm. Select the number of duplicates determined above, In the execution step, when the own device is newly selected by the selection step to the computer on which the arbitrary data is newly arranged, the arbitrary data is written to the computer, and the computer newly arranges the arbitrary data by the selection step. A data processing method characterized in that the data is deleted when the data is no longer selected.
<figref num="1">It is explanatory drawing which shows the system configuration of the server apparatus which concerns on this Embodiment.</figref><figref num="2">It is a block diagram which shows the hardware configuration of the server apparatus which concerns on this embodiment.</figref><figref num="3">It is a block diagram which shows the functional structure of the cooperation processing part.</figref><figref num="4">It is explanatory drawing which shows the data update processing in a server apparatus.</figref><figref num="5">It is a flowchart which shows the procedure of the arrangement server selection process.</figref><figref num="6-1">It is a flowchart which shows the procedure of data update processing (in the case of the server device which received the update request from the application execution part of own device).</figref><figref num="6-2">It is a flowchart which shows the procedure of data update processing (in the case of the server device which received the update request from the cooperation processing part of another server device).</figref><figref num="7">It is explanatory drawing which shows the data reference processing in a server apparatus.</figref><figref num="8-1">It is a flowchart which shows the procedure of data reference processing (in the case of the server device which received the reference request from the application execution part of own device).</figref><figref num="8-2">It is a flowchart which shows the procedure of data reference processing (in the case of the server device which received the reference request from the cooperation processing part of another server device).</figref><figref num="9">It is a flowchart which shows the data rearrangement processing at the time of changing the number of replicas.</figref><figref num="10">It is a flowchart which shows the data rearrangement processing when the number of replicas increases.</figref><figref num="11">It is a flowchart which shows the data rearrangement processing when the number of replicas decreases.</figref><figref num="12">It is explanatory drawing which shows the latency calculation procedure at the time of data update.</figref><figref num="13">It is explanatory drawing which shows the latency calculation procedure at the time of data reference.</figref>
Code description
100 server unit 101 App Execution Department 102 Coordination processing unit 103 Memory 110 load balancer 120 requests 301 Acquisition Department 302 Selection 303 Transmitter / receiver 304 Setting section 305 Judgment Department 306 Execution section 307 Decision Department
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 | Cited during |
|---|---|---|---|
| JP2014115929A | Cited by | Japan | Search report |
| US10691716B2 | Cited by | United States of America | Applicant |
| US10956246B1 | Cited by | United States of America | Applicant |
| JP5765416B2 | Cited by | Japan | Search report |
| WO2014119269A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| JPWO2014119269A1 | Cited by | Japan | Search report |
| JP2012043358A | Cited by | Japan | Examiner |
| US11070600B1 | Cited by | United States of America | Applicant |
| US9715521B2 | Cited by | United States of America | Applicant |
| WO2012121316A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| JP2016507079A | Cited by | Japan | Examiner |
| CN105765575A | Cited by | China | Search report |
| US10855754B1 | Cited by | United States of America | Applicant |
| WO2015166741A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| JP2015512551A | Cited by | Japan | Search report |
| JP2013045181A | Cited by | Japan | Examiner |
| US10635644B2 | Cited by | United States of America | Applicant |
| US9830324B2 | Cited by | United States of America | Applicant |
| US9846553B2 | Cited by | United States of America | Applicant |
| US11509700B2 | Cited by | United States of America | Applicant |
| US10691716B2 | Cited by | United States of America | Applicant |
| WO2012121316A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| US12375556B2 | Cited by | United States of America | Applicant |
| US11621999B2 | Cited by | United States of America | Applicant |
| JP2017501515A | Cited by | Japan | Search report |
| JP2012221419A | Cited by | Japan | Examiner |
| JP2013030035A | Cited by | Japan | Examiner |
| JP2017501515A | Cited by | Japan | Search report |
| US10768830B1 | Cited by | United States of America | Applicant |
| US10248556B2 | Cited by | United States of America | Applicant |
| US10474654B2 | Cited by | United States of America | Applicant |
| US9628438B2 | Cited by | United States of America | Applicant |
| US9985829B2 | Cited by | United States of America | Applicant |
| US9774582B2 | Cited by | United States of America | Applicant |
| US9342574B2 | Cited by | United States of America | Applicant |
| US11675501B2 | Cited by | United States of America | Applicant |
| US8972365B2 | Cited by | United States of America | Applicant |
| WO2014010023A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| US10467105B2 | Cited by | United States of America | Applicant |
| US11075984B1 | Cited by | United States of America | Applicant |
| US10798140B1 | Cited by | United States of America | Applicant |
| JP2017501515A | Cited by | Japan | Search report |
| US9934242B2 | Cited by | United States of America | Applicant |
| JP2015212855A | Cited by | Japan | Search report |
| JP2001067377A | Cites | Japan | Examiner |
| WO2006059476A1 | Cites | World Intellectual Property Organization (WIPO) | Examiner |
| JP2007065714A | Cites | Japan | Examiner |
| JP2007133503A | Cites | Japan | Examiner |
4 members in 2 offices
Priority claims2
| Document | Office | Kind | Date |
|---|---|---|---|
| 2008319530 | Japan | A | |
| JP20080319530 | – | – | – |
Members4
| Document | Office | Kind | |
|---|---|---|---|
| US2010153337A1 | United States of America | A1 | |
| JP2010146067AThis record | Japan | A | |
| US8577838B2 | United States of America | B2 | |
| JP5396848B2 | Japan | B2 |
7 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 | |
| Report on retrievalJAPANESE INTERMEDIATE CODE: A971007A977 | A977 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 2010146067
- Publication, DOCDB
- 2010146067
- Publication, EPODOC
- JP2010146067
- Application
- 319530
- Application, DOCDB
- 2008319530
- Application, EPODOC
- JP20080319530
Titles2
- Japanese
- データ処理プログラム、サーバ装置およびデータ処理方法
- English
- Data processing program, server device and data processing method
Classification
- CPC, 2
- G06F9/505
- G06F16/27
- IPC, 1
- G06F12 00