Program execution load distribution method in computer system
Abstract
(57) By assigning a method automatically to another computer of low load, and performing it for every class, at the time of execution of an application program, the summary purpose book invention distributes load, and eliminates the method assigned at the time of an end from another computer. Composition Two or more computers are connected through the communication medium 20. When a method is performed during execution of the application program 22 on OS2110 of a certain computer 211, the load sharing means 23, If the class to which the method belongs is not assigned to all the computer, either, the computer of low load will be searched, the class will be assigned, and all the methods belonging to the class of will be transmitted. And -- about the method which belongs to the class concerned when the class to which the method to perform belongs is assigned to other computers, it is execution -- being concerned -- others -- a computer is requested and the class to which the method belongs performs by a self-computer about the method which belongs to a quota たら当 this class at a self-computer. If the application program 22 is completed, the load sharing means 23 will eliminate all the methods belonging to the class assigned to the other computers concerned from other computers.
Term
Term ended
Projected expiry passed 7 June 2013, 13.3 years ago.
- Priority and filed
- Published
- Projected expiry
- Today
5 claims: 3 independent, 2 dependent
- 1[Claims] 1. A method of distributing the load of executing a program executed by a specific computer to a plurality of computers including the specific computer in a computer system having two or more computers connected via a communication medium. And A step of providing a table in which a computer to which a processing unit, which is a unit of processing constituting the program, is assigned to the specific computer is registered in association with the assigned processing unit. A step in which the specific computer sequentially executes the processing unit, and When the specific computer executes the processing unit, a step of referring to the table to determine whether or not a computer to which the processing unit is assigned exists, and When it is determined that the computer to which the processing unit is assigned does not exist, the processing unit is assigned to one of the plurality of computers, and the assigned computer is associated with the assigned processing unit to the table. A step of registering and transferring the processing unit to the assigned computer from the specific computer and executing the process. If it is determined that there is a computer to which the processing unit is assigned, the step of causing the computer to execute the processing unit and When the specific computer finishes the execution of the program, all the processing units transferred from the specific computer to the other computer are deleted from the other computer, and the table is initialized. A program execution load distribution method in a computer system characterized by having. 【特許請求の範囲】 【請求項1】通信媒体を介して接続された2以上の計算機を有する計算機システムにおいて、特定の計算機で実行されるプログラムの実行の負荷を、前記特定の計算機を含む複数の計算機に分散させる方法であって、 前記特定の計算機に、前記プログラムを構成する処理の単位である処理単位を割り当てた計算機を、割り当てた処理単位に関連付けて登録するテ-ブルを設けるステップと、 前記特定の計算機が、前記処理単位を順次実行するステップと、 前記特定の計算機が前記処理単位を実行する際に、前記テ-ブルを参照して、当該処理単位が割り当てられた計算機が存在するか否かを判断するステップと、 処理単位が割り当てられた計算機が存在しないと判断された場合に、当該処理単位を前記複数の計算機の内の一の計算機に割り当て、割り当てた計算機を割り当てた処理単位に関連付けて前記テ-ブルに登録し、割り当てた計算機に当該処理単位を前記特定の計算機より転送して実行させるステップと、 処理単位が割り当てられた計算機が存在すると判断された場合には、当該計算機に当該処理単位を実行させるステップと、 前記特定の計算機が、前記プログラムの実行を終了した際に、前記特定の計算機より他の計算機に転送した処理単位をすべて、当該他の計算機より消去し、前記テ-ブルを初期化するステップとを有することを特徴とする計算機システムにおけるプログラム実行負荷分散方法。
- 3A method of distributing the load of executing a program executed by a specific computer to a plurality of computers including the specific computer in a computer system having two or more computers connected via a communication medium. And A step of providing a table in which a computer to which a class, which is a unit of a processing group constituting the program, is assigned to the specific computer is registered in association with the assigned processing unit. A step in which the specific computer sequentially executes a method that is a process belonging to the class. When the specific computer executes an arbitrary method, a step of referring to the table to determine whether or not there is a computer to which a class to which the arbitrary method belongs exists. When it is determined that there is no computer to which the class to which the arbitrary method belongs is determined, the class is assigned to one of the plurality of computers, and the assigned computer is associated with the assigned class. -The step of registering in the bull, transferring all the methods belonging to the class to the assigned computer from the specific computer, and executing the arbitrary method. If it is determined that there is a computer to which the class to which the arbitrary method belongs is determined, the step of causing the computer to execute the arbitrary method and When the specific computer finishes executing the program, all the methods transferred from the specific computer to another computer are deleted from the other computer, and the table is initialized. A program execution load distribution method in a computer system characterized by having. 【請求項3】通信媒体を介して接続された2以上の計算機を有する計算機システムにおいて、特定の計算機で実行されるプログラムの実行の負荷を、前記特定の計算機を含む複数の計算機に分散させる方法であって、 前記特定の計算機に、前記プログラムを構成する処理群の単位であるクラスを割り当てた計算機を、割り当てた処理単位に関連付けて登録するテ-ブルを設けるステップと、 前記特定の計算機が、前記クラスに属する処理であるメソッドを順次実行するステップと、 前記特定の計算機が任意の前記メソッドを実行する際に、前記テ-ブルを参照して、当該任意のメソッドの属するクラスが割り当てられた計算機が存在するか否かを判断するステップと、 前記任意のメソッドの属するクラスが割り当てられた計算機が存在しないと判断された場合に、当該クラスを前記複数の計算機の内の一の計算機に割り当て、割り当てた計算機を割り当てたクラスに関連付けて前記テ-ブルに登録し、割り当てた計算機に当該クラスに属する全てのメソッドを前記特定の計算機より転送し、前記任意のメソッドを実行させるステップと、 前記任意のメソッドの属するクラスが割り当てられた計算機が存在すると判断された場合には、当該計算機に前記任意のメソッドを実行させるステップと、 前記特定の計算機が、前記プログラムの実行を終了した際に、前記特定の計算機より他の計算機に転送したメソッドを、当該他の計算機より全て消去し、前記テ-ブルを初期化するステップとを有することを特徴とする計算機システムにおけるプログラム実行負荷分散方法。
- 5A computer system including a communication medium and a plurality of computers connected to each other via the communication medium. The calculator The storage unit that stores the methods of the program to be executed, the computer to which the class that is the unit of processing that constitutes the program to be executed is assigned, and the load amount of execution of all the methods belonging to the class are assigned to the assigned processing unit. The load management unit that associates and stores, the method that executes the method that executes the method requested to be executed, the method execution means that reports the load required to execute the method to the requester, and the request for execution are requested. A method proxy execution means that requests the method execution means to execute the method and relays the report of the load amount to the requester, and a method proxy execution means that executes the method requested to be executed by the method of another computer. The method execution request means for relaying the report of the load amount to the requester, the class storage means for storing the transferred method in the storage unit, and all the methods belonging to the class requested to be transferred are described above. A class transfer means that transfers from the storage unit to the class storage means of another computer, When processing an arbitrary method, the computer to which the class to which the arbitrary method belongs is checked from the stored contents of the load management unit, and the computer to which the class to which the arbitrary method belongs does not exist. In this case, the computer having the smallest cumulative load of executing the method stored in association with the load management unit is determined as the computer to which the class to which the arbitrary method belongs is determined, and the class transfer means is determined. The computer is requested to transfer the class to which the arbitrary method belongs, and when the class to which the arbitrary method belongs is assigned to the own computer, the method executing means is requested to execute the arbitrary method, and the arbitrary method is executed. When the class to which the method belongs is assigned to another computer, the method execution request means is requested to request the other computer to execute the arbitrary method, and the load management unit is informed of the load charge. According to the report, the execution selection means that stores the load amount required to execute the methods belonging to each class in the load management unit in association with the computer to which the class unit is assigned, and Class erasing means that erases the method belonging to the class requested to be erased from the storage unit and initializes the stored contents of the load management unit in response to the request, erasing the class requested to be erased, and the above. A class erasing requesting means that requests the initialization of the stored contents of the load management unit to the class erasing means of another computer, and a class erasing requesting means. When the execution of the program is completed, the class erasing means is requested to initialize the stored contents of the load management unit according to the stored contents of the load management unit, and the class erasing requesting means is used to erase the class. A computer system characterized by having a program termination means for requesting a request to a device class erasing means. 【請求項5】通信媒体と、当該通信媒体を介して互いに接続された複数の計算機を有する計算機システムであって、 前記計算機は、 実行するプログラムのメソッドを記憶した記憶部と、実行するプログラムを構成する処理の単位であるクラスを割り当てた計算機と、当該クラスに属する全てのメソッドの実行の負荷量とを、割り当てた処理単位に関連付けて記憶する負荷管理部と、実行を依頼されたメソッドを実行するメソッドを実行し、当該メソッドの実行に要した負荷量の報告を依頼元に行うメソッド実行手段と、実行の依頼を依頼されたメソッドの実行を前記メソッド実行手段に依頼し、前記負荷量の報告を依頼元に中継するメソッド代理実行手段と、実行の依頼を依頼されたメソッドの実行を他の計算機のメソッドの代理実行手段に依頼し、前記負荷量の報告を依頼元に中継するメソッド実行依頼手段と、転送されたメソッドを前記記憶部に記憶するクラス格納手段と、転送を依頼されたクラスに属する全てのメソッドを前記記憶部より他の計算機のクラス格納手段に転送するクラス転送手段と、 任意の前記メソッドを処理する際に、前記負荷管理部の記憶内容より、当該任意のメソッドの属するクラスが割り当てられた計算機をチェックし、前記任意のメソッドの属するクラスが割り当てられた計算機が存在しない場合に、前記負荷管理部に対応付けられて記憶されているメソッドの実行の負荷量の累積が最も小さい計算機を前記任意のメソッドの属するクラスを割り当てる計算機として決定し、前記クラス転送手段に決定した計算機に前記任意のメソッドの属するクラスの転送を依頼し、前記任意のメソッドの属するクラスが自計算機に割り当てられている場合に、前記メソッド実行手段に前記任意のメソッドの実行を依頼し、前記任意のメソッドの属するクラスが他計算機に割り当てられている場合に、当該他の計算機への前記任意のメソッドの実行の依頼を前記メソッド実行依頼手段に依頼し、前記負荷管理部に、前記負荷料の報告に応じて、それまでに各クラスに属するメソッドの実行に要した負荷量を、当該クラス単位を割り当てる計算機に対応付けて前記負荷管理部に記憶する実行選択手段と、 消去を依頼されたクラスに属するメソッドを前記記憶部より消去し、依頼に応じて前記負荷管理部の記憶内容を初期化するクラス消去手段と、消去の依頼を依頼されたクラスの消去と、前記負荷管理部の記憶内容の初期化を他の計算機の前記クラス消去手段に依頼するクラス消去依頼手段と、 前記プログラムの実行が終了した場合に、前記負荷管理部の記憶内容に応じて、クラス消去手段に前記負荷管理部の記憶内容の初期化を依頼し、前記クラス消去依頼手段にクラスの消去の他装置のクラス消去手段への依頼を依頼するプログラム終了手段とを有することを特徴とする計算機システム。
Independent claims3
203 paragraphs, as filed
Description: TECHNICAL FIELD [Detailed description of the invention]
【0001】
[Industrial application field]
The present invention relates to a distributed computer system including a plurality of computers connected by a communication medium, and more particularly to a technique for distributing the load of application program execution in the distributed computer system.
【0002】
[Conventional technology]
Conventionally, as a technique for distributing the load of application program execution in a distributed computer system in which a plurality of computers are connected by a communication medium, an application program is prepared in advance in a plurality of computers in the system, and these multiple computers are executed when the application program is executed. There is known a technique for identifying a computer that seems to have the highest execution efficiency and requesting the computer to execute an application program in a batch.
【0003】
In addition, the techniques described in JP-A-3-42765 and JP-A-2-285460 are known.
【0004】
In this technique described in Japanese Patent Application Laid-Open No. 2-285460, when executing an application program, a computer that seems to have the highest execution efficiency among a plurality of computers in the system is identified, and at that time, the computer is used. The application programs are collectively transferred and executed.
【0005】
Further, in the technique described in Japanese Patent Application Laid-Open No. 3-42765, the application program is divided into predetermined processing units in advance, and each processing unit is provided in advance in a plurality of computers in the system. Then, when the corresponding processing unit is called during the execution of the application program, the computer that seems to have the highest execution efficiency is specified in consideration of the execution cost and the like, and the processing unit is executed for that computer. I'm asking.
【0006】
[Problems to be Solved by the Invention]
These conventional methods have the following problems.
【0007】
That is, in the technique of requesting the execution of the application program to another computer, this application program must be prepared in advance in the other computer to which the execution is requested. Therefore, since each application program for sufficiently distributing the load must be prepared for each computer constituting the distributed computer system, the amount of resources of each computer required for load distribution becomes excessive.
【0008】
Further, in the technique described in Japanese Patent Application Laid-Open No. 2-285460, when the scale of the application program becomes large, the load cannot be sufficiently distributed to each computer.
【0009】
Further, in the technique described in Japanese Patent Application Laid-Open No. 3-42765, the processing unit to which the application program requests execution must be provided in advance in the other computer, which increases the resource amount of each computer required for load distribution. .. Furthermore, before requesting execution of a processing unit, the processing execution cost in each computer for each processing unit must be specified, but the contents of the table that manages this execution cost differ depending on the computer that executes the application program. It has to be a value. Therefore, when changing the calculator that executes the application program from one calculator to another, the contents of this table must be recreated. Therefore, it is not always suitable for load distribution of an application program that is desired to be used by an arbitrary computer.
【0010】
Therefore, an object of the present invention is to provide a distributed computer system capable of sufficiently distributing the load of executing an application program without increasing the amount of resources of each computer required for load distribution.
【0011】
[Means for solving problems]
In order to achieve the above object, the present invention has a plurality of computers including the specific computer, in which the load of executing a program executed by a specific computer is applied in a computer system having two or more computers connected via a communication medium. A method of distributing to computers, in which a step of providing a table in which a computer to which a processing unit, which is a unit of processing constituting the program, is assigned to the specific computer is registered in association with the assigned processing unit is provided. There is a step in which the specific computer sequentially executes the processing unit, and a computer to which the processing unit is assigned with reference to the table when the specific computer executes the processing unit. The step of determining whether or not, and when it is determined that the computer to which the processing unit is assigned does not exist, the processing unit is assigned to one of the plurality of computers, and the assigned computer is assigned. When it is determined that there is a step of registering in the table in association with a unit and transferring the processing unit from the specific computer to the assigned computer for execution, and a computer to which the processing unit is assigned, if it is determined that there is a computer to which the processing unit is assigned. The step of causing the computer to execute the processing unit and all the processing units transferred from the specific computer to another computer when the specific computer finishes executing the program are deleted from the other computer. The present invention provides a method for distributing a program execution load in a computer system, which comprises a step of initializing the table.
【0012】
Here, the table is associated with a registered computer, and the load amount required for execution of the processing unit assigned to the computer is set as a table that can be registered, and the processing unit is executed. At the end, the load amount already registered associated with the computer to which the processing unit is assigned is updated to the value obtained by adding the load amount required to execute the method whose execution has been completed to the load amount. Further, the processing unit is assigned to the computer having the step further, and the processing unit is assigned to the computer having the smallest cumulative execution load of the processing unit registered in association with the table. You may do it.
【0013】
[Action]
According to the program execution load distribution method according to the present invention, since the processing unit is automatically assigned to each computer during the execution of the program, it is not necessary to prepare the processing unit in another computer in advance. Therefore, the program creator and the user are not particularly aware of the distributed processing and the computer that actually performs the processing. In addition, all processing units assigned to other computers are deleted when the application program ends, eliminating waste of resources.
【0014】
In addition, when the program is executed, the computer to be assigned is selected based only on the execution load of the processing unit, so the load is distributed among the computers that make up the distributed computer system, and the program can be operated in the same way on any computer in the system. Can be done.
【0015】
[Example]
Hereinafter, an embodiment of the distributed computer system according to the present invention will be described.
【0016】
Figure 1 shows the configuration of the distributed computer system according to this embodiment.
【0017】
As shown in the figure, the computers 211, 212, and 213 are connected via the communication medium 20 to form a distributed computer system. Each computer 211, 212, 213 is a standard computer including a processor, a storage device, a communication control device, and the like. Each computer 211, 212, 213 is equipped with an operating system (OS) 2110, 2120, 2130, and executes an application program on the OS. For example, in computer 211, application program 22 is executed on OS 2110. Each OS 2110, 2120, 2130 is provided with load distribution means 23, 232, and 233 to distribute the load of executing the application program, respectively. For example, the load distribution means 23 transfers the processing unit of the application program 22 to the load distribution means 232 and 233 of the other computers 212 and 213 and requests execution. The load distribution means 232 and 233 of the other computers 212 and 213 execute the processing unit requested to be executed, and return the result to the request source load distribution means 23.
【0018】
Next, the application program 22 will be described.
【0019】
The application program 22 is a program written in an object-oriented programming language. In this embodiment, it is described in a programming language widely known as "C ++".
【0020】
Figure 2 shows the application program 22.
【0021】
As shown, the application program 220 consists of class definitions 51, 52, 53 and a main program 54.
【0022】
Here, the class will be described with reference to FIG.
【0023】
As shown in FIG. 3, the class defined as class X10 consists of a data unit 11 which is a set of data definitions and a method unit 12 which is a set of methods which are operations for manipulating data. The data unit 11 is composed of class-specific data definitions 111 to 119, and the method unit 12 is composed of methods 121 to 129, which are processing units for manipulating data.
【0024】
Objects 13 and 14 in the figure are variables whose type is declared to be class X10. Since objects 13 and 14 represent an instance of class 10, the method and data definitions are common, but the data itself has its own unique one.
【0025】
Now, returning to Fig. 2, in this program, first, classes A, B, and C are defined by class definitions 51, 52, and 53, respectively. In the main program 54, object declarations 541, 542, and 543 declare class A object c1, class B object c2, and class C object c3, respectively. Next, 544 specifies that the method func1 of the class A to which the object c1 belongs is executed for the object c1. Similarly, 545 specifies that object c2 executes the method func1 of class B to which object c2 belongs, and 546 executes the method func2 of class A to which object c1 belongs to object c1. Is specified.
【0026】
The OS2110 executes such an application program 22 sequentially. At this time, the load distribution means 23, 232, and 233 realize the load distribution as described above.
【0027】
FIG. 4 shows the functional configuration of the load distribution means 23. The load distribution means 232 and 233 also have the same configuration.
【0028】
In FIG. 4, the communication control means 30 is connected to the communication medium 20 and communicates between computers. The storage unit 31 stores application programs and objects executed on the self-computer, method units of classes assigned to the self-computer by other computers, and object data sent from other computers to the self-computer. .. The class management unit 32 stores the information of the class assigned to the self-calculator. The method execution means 33 executes the method in the storage unit 31 based on the information in the class management unit 32. The method execution request means 34 requests the designated computer to execute the method via the communication control means 30, and obtains the execution result. The class transfer means 35 transfers the method unit of the class to the designated computer via the communication control means 30. The class erasure request means 36 requests the designated computer via the communication control means 30 to erase the method unit of the designated class. The method proxy execution means 37 calls the method execution means 33 to realize the execution of the method requested by the method execution request means 34 via the communication control means 30. The class storage means 38 stores the method unit transferred by the class transfer means 35 via the communication control means 30 in the storage unit 31. The class erasing means 39 erases the method unit of the class designated by the class erasing requesting means 36 from the storage unit 31 via the communication control means 30. The load management unit 40 manages load information regarding the class to which the method execution request means 34 and the method execution means 33 execute. When a method in a certain class is called for the first time from an application program, the execution selection means 41 specifies a computer having a low load based on the information of the load management unit 40, calls the class transfer means 35, and transfers the method part of the class. , And when the method is called from the application program when the class to which the method belongs has already been transferred by the class transfer means 35, specify which computer the class to which the method belongs to is based on the information of the load management unit 40. , When the designated calculator is its own calculator, the method execution means 33 is called, and when it is another calculator, the method execution request means 34 is called, and the means for updating the information of the load management unit 40 is included. The program termination means 42 specifies a class and a calculator based on the information of the load management unit 40, calls the class deletion request means 36, and deletes the class from the load management unit 40, which is stored in the load management unit 40. Repeat until all are erased.
【0029】
Hereinafter, details of each part of the load distribution means shown in FIG. 4 will be described.
【0030】
First, as shown in FIG. 5, the load management unit 40 has a class name 401, a computer name 402 that executes a method of the class 401, a processing time 403 that is the total execution time of all the methods of the class 401, and a computer. It consists of a table that stores an identification number 404 for identifying class 401.
【0031】
Further, as shown in FIG. 6, the class management unit 32 includes a table that stores the class name 321 and the identification number 322 of the class 321 in the computer.
【0032】
Next, the execution selection means 41 is a means that is called with the method name to be executed and its class as arguments and performs processing according to the flowchart shown in FIG.
【0033】
That is, the execution selection means 41 first searches for a class from the class name 401 of the load management unit 40 (step 10010). Check if the class name can be searched (10020), and if it can be searched, obtain the computer name 402 corresponding to the class name and identify the computer (10070).
【0034】
If the search cannot be performed, the total processing time 403 is calculated for each computer name 402 in the load management unit 40, the computer with the smallest total processing time is specified (10030), and whether it is a self-calculator (10040). If it is not a self-calculator, call class transfer means 35 (10050) and request transfer of the method part of the class to which the method to be executed belongs. Then, the class name of the transferred method unit and the transferred computer name are registered in the class name 401 and the computer name 402 of the load management unit 40, and the processing time 403 is initialized (10060).
【0035】
If the machine that executes the method can be identified, it determines whether it is a self-calculator or another computer (10080), and if it is a self-calculator, calls the method execution means 33 to execute the method on the self-calculator (10090), otherwise the method is executed. Call request means 34 to execute the method on another computer (10100).
【0036】
After that, the method execution time is returned from the method execution means 33 or the method execution request means 34, so add it to the processing time 403 of the load management unit 40 corresponding to the method class (10110), and call the processing. Undo.
【0037】
Next, the class transfer means 35 is a means that is called with the partner computer name and the class name as arguments and executes processing according to the flowchart shown in FIG.
【0038】
That is, the class transfer means 35 first transfers the method unit of the class from the storage unit 31 to the class storage means 38 of the other computer (15010), returns the returned identification number as a return value (15020), and processes. Is returned to the caller.
【0039】
Next, the class storage means 38 is a means for performing processing according to the flow chart shown in FIG.
【0040】
That is, the class storage means 38 stores the method unit transferred from another computer in the storage unit 31 (13010), registers the class name of the method unit in the class name 321 of the class management unit 32, and corresponds to the method unit. The identification number 322 is returned as a return value (13020), and the process is returned to the caller.
【0041】
Next, the method execution means 33 is a means that is called with an argument of the method name and an identification number for the computer to identify the class, and performs processing according to the flowchart shown in FIG.
【0042】
The method execution means 33 first selects the corresponding class name 321 from the identification number 322 of the class management unit 32 for the method class, and specifies the class (11010). Then, the method corresponding to the method name passed as an argument of the specified class stored in the storage unit 31 is executed using the data of the object stored in the storage unit 31 (11020). Then, the execution result and the execution time are used as return values (11030), and the process is returned to the caller.
【0043】
Next, the method execution request means 34 is a means that is called with an argument of the counterpart computer name, the method name, and the identification number that the computer identifies the class, and executes the process according to the flowchart shown in FIG.
【0044】
That is, the method execution request means 34 first calls the method proxy execution means 37 of the partner computer via the communication control means 30 with the method name and the class identification number as arguments (14010). At this time, the data (object data) required for executing the method is read from the storage unit 31 and passed to the method proxy execution means 37. Then, the method proxy execution means 37 of the remote computer returns the returned execution result and execution time as a return value (14020), and returns the process to the caller.
【0045】
Next, the method proxy execution means 37 is a means that is called with the method name and the number that the other computer identifies the class as arguments, and performs processing according to the flowchart shown in FIG.
【0046】
First, the method proxy execution means 37 stores the data (object data) necessary for executing the passed method in the storage unit 31, and calls the method execution means 33 with the identification number and the method name as arguments (12010). .. Then, when the execution result and the execution time are returned from the method execution means 33, this is used as the return value (12020) and the process is returned to the caller. At this time, the data (object data) necessary for executing the method stored in the storage unit 31 is deleted.
【0047】
Next, the program termination means is a means for performing the operation according to the flowchart of FIG.
【0048】
When the program termination means 42 is called, first, the class name 401, the computer name 402, and the identification number 404 of the load management unit 40 are obtained (16010), and while the class name exists (16020), the computer name is self-identified. Check if it is a computer (16030), if it is a self-computer, call the class erasure means with the identification number as an argument (16040), otherwise call the class erasure request means 36 with the computer name and identification number as arguments (16050), Delete the class name 401, computer name 402, processing time 403, and identification number 404 of the load management unit 40 (16050). When the class name disappears, the process is returned to the caller.
【0049】
Next, the class erasing means 39 is called with the identification number as an argument, and is a means for performing processing according to the flowchart shown in FIG.
【0050】
That is, if the class erasing means 39 searches for the identification number from the identification number 322 of the class management unit 32, identifies the class from the corresponding class name 321 (17010), and the class is assigned by another computer. For example, the method part of the class that has been transferred and stored in the storage part 31 is deleted (17020). In addition, the class and identification number are deleted from the class name 321, number 322 of the class management unit 32 (17030), and the process is returned to the caller.
【0051】
Next, the class erasure request means 36 is a means that is called with the partner computer name and the identification number as arguments and performs processing according to the flowchart shown in FIG.
【0052】
That is, the class erasure request means 36 first calls the class erasure means 39 of the other computer (18010) via the communication control means 30, and returns the process to the caller.
【0053】
Hereinafter, the execution operation of the application program of the distributed computer system according to this embodiment will be described.
【0054】
First, FIG. 16 shows a processing procedure performed by OS2110 for executing the application program 22.
【0055】
As shown in the figure, the OS 2110 first sequentially executes the main program 54 of the application program 22 stored in the storage unit 31 (19040). If you experience the execution of the method in the middle of processing (19010), the method name and the class name as an argument the actual load balancing means 23 calls the row selection means 41 (19020). Also, when the main program ends (19030), the program termination means 42 is called (19050).
【0056】
When the application program 22 shown in FIG. 2 is executed by such an operation of OS 2110, the operation of each part is as follows.
【0057】
First, the application program 22 main program 54 is executed sequentially. Figure 17 shows the state of the distributed computer system immediately before the start of execution of application program 22.
【0058】
First, object declarations 541, 542, and 543 declare class A object c1, class B object c2, and class C object c3, respectively. Next, since 544 is a method, execution selection means 41 is called. The execution selection means 41 determines the computer to which the class A is assigned because the method 544 is a method of the class A and the class A is not in the class name 6411 of the load management unit 641.
【0059】
Since there is no information in the load management unit 641, it is randomly determined and assigned to the computer 611. Since the computer 611 is a computer running the application program 22, there is no need to transfer the methods of the class.
【0060】
The class A and the computer 611 are stored in the class name 6411 and the computer name 6412 of the load management unit 641 and the processing time 6413 is initialized. Then, the execution selection means 41 calls the method execution means 33. Method execution means 33 executes method 544 and returns the execution result and processing time. The execution selection means 41 adds the returned processing time to the one corresponding to the class A of the processing time 6413 of the load management unit 641. The execution selection means returns the method execution result to the application program 22. Figure 18 shows the state of the distributed computer system at this point.
【0061】
Then 545 in Figure 2 is executed. Since 545 is a method, execution selection means 41 is called. The execution selection means 41 determines the computer to which the class B is assigned because the method 545 is a method of the class B and the class B is not in the class name 6411 of the load management unit 641. Since the processing time of computers 612 and 613 is 0 from the computer name 6412 and processing time 6413 of the load management unit 641, it is randomly determined from these two and assigned to the computer 613. Since the computer 613 is not the computer 611 executing the application program 22, the class transfer means 35 is called. The class transfer means 35 transfers the method unit of the class B to the class storage means 38 of the computer 613 via the communication control means 30. The class storage means 38 stores the transferred method unit in the storage unit 653, stores the class B in the class name 6631 of the class management unit 663, determines the identification number of the class B, stores it in the identification number 6632, and returns the value. The execution selection means 41 stores the returned identification number, class B, and computer 613 in the identification number 6414, class name 6411, and computer name 6412 of the load management unit 641 and initializes the processing time 6413.
【0062】
Then, the execution selection means 41 calls the method execution request means 34. The method execution request means 34 calls the method proxy execution means 37 of the computer 613 via the communication control means 30. Method proxy execution means 37 specifies method 545 and calls the method execution means. The method execution means executes method 545 and returns the execution result and processing time. The execution selection means 41 adds the returned processing time to the class B corresponding to the processing time 6413 of the load management unit 641. The execution selection means returns the method execution result to the application program 22. Figure 19 shows the state of the distributed computer system at this point.
【0063】
Then 546 in Figure 2 is executed. Since 546 is a method, execution selection means 41 is called. Since the method 546 is a class A method, the execution selection means 41 searches the class name 6411 and the computer name 6412 of the load management unit 641 and knows that the method part of the class A is assigned to the computer 611. Since the computer 611 is a computer 611 (self-computer) executing the application program 22, the execution selection means 41 calls the method execution means. Method execution means 33 executes method 546 and returns the execution result and processing time. The execution selection means 41 adds the returned processing time to the one corresponding to the class A of the processing time 6413 of the load management unit 641. The execution selection means returns the method execution result to the application program 22. Figure 20 shows the state of the distributed computer system at this point.
【0064】
Next, 547 in Figure 2 is executed. Since 547 is a method, execution selection means 41 is called. The execution selection means 41 determines the computer to which the class C is assigned because the method 547 is a method of the class C and the class C is not in the class name 6411 of the load management unit 641. Since the processing time of the computer 612 is 0, which is the lowest value from the computer name 6412 and the processing time 6413 of the load management unit 641, it is assigned to the computer 612. Since the computer 612 is not the computer 611 executing the application program 22, the class transfer means 35 is called. The class transfer means 35 transfers the method unit of the class C to the class storage means 38 of the computer 612 via the communication control means 30. The class storage means 38 stores the transferred method unit in the storage unit 653, stores the class C in the class name 6631 of the class management unit 663, determines the identification number of the class C, stores it in the identification number 6632, and returns the value. The execution selection means 41 stores the returned identification number, class C, and computer 612 in the identification number 6414, class name 6411, and computer name 6412 of the load management unit 641 and initializes the processing time 6413. Then, the execution selection means 41 calls the method execution request means 34. The method execution request means 34 calls the method proxy execution means 37 of the computer 612 via the communication control means 30. Method proxy execution means 37 specifies method 545 and calls the method execution means. The method execution means executes method 545 and returns the execution result and processing time. The execution selection means 41 adds the returned processing time to the class B corresponding to the processing time 6413 of the load management unit 641. The execution selection means returns the method execution result to the application program 22. Figure 21 shows the state of the distributed computer system at this point.
【0065】
Next, 548 in Figure 2 is executed. Since 548 is a method, execution selection means 41 is called. Since the method 548 is a class B method, the execution selection means 41 searches the class name 6411 and the computer name 6412 of the load management unit 641 and knows that the method part of the class B is assigned to the computer 613. Since the computer 613 is not the computer 611 executing the application program 22, the execution selection means 41 calls the method execution request means 34. The method execution request means 34 calls the method proxy execution means 37 of the computer 613 via the communication control means 30. Method proxy execution means 37 specifies method 545 and calls the method execution means. The method execution means executes method 545 and returns the execution result and processing time. The execution selection means 41 adds the returned processing time to the class B corresponding to the processing time 6413 of the load management unit 641. The execution selection means returns the method execution result to the application program 22. Figure 22 shows the state of the distributed computer system at this point.
【0066】
Now that the execution of the main program 54 of the application program 22 of FIG. 2 is completed, the program termination means 42 is called. The program termination means 42 searches the class name 6411 and the computer name 6412 of the load management unit 641 and finds out that the method unit of the class A is assigned to the computer 611. Since the computer 611 is a computer 611 (self-computer) executing the application program 22, the class erasing means 39 is called by designating the identification number 6414. The class erasing means 39 deletes the class name 321 and the identification number 322 of the class management unit 32 corresponding to the specified identification number.
【0067】
The program termination means 42 deletes the class name 6411 of the load management unit 641 and the corresponding computer name 6412, processing time 6413, and identification number 6414, continues the search, and finds out that the method part of the class B is assigned to the computer 613. Since the computer 613 is not the computer 611 executing the application program 22, the class deletion request means 36 is called by specifying the computer name 6412 and the identification number 6414.
【0068】
The class erasing request means 36 calls the class erasing means 39 of the computer 613 via the communication control means 30. The class erasing means 39 deletes the class name 321 and the identification number 322 of the class management unit 32 corresponding to the specified identification number. The program termination means 42 deletes the class name 6411 of the load management unit 641 and the corresponding computer name 6412, processing time 6413, and identification number 6414, continues the search, and finds out that the method part of the class C is assigned to the computer 612.
【0069】
Since the computer 612 is not the computer 611 executing the application program 22, the class deletion request means 36 is called by specifying the computer name 6412 and the identification number 6414. The class erasing request means 36 calls the class erasing means 39 of the computer 612 via the communication control means 30. The class erasing means 39 deletes the method of the class corresponding to the specified identification number from the storage unit 31, and deletes the class name 321 and the identification number 322 of the corresponding class management unit 32. The program termination means 42 deletes the class name 6411 of the load management unit 641 and the corresponding computer name 6412, processing time 6413, and identification number 6414, continues the search, and finds out that the class is no longer stored in the class name 6411 and performs processing. Finish. Figure 23 shows the state of the distributed computer system at this point, and it can be seen that all the methods of the application program 22 have been deleted from the distributed computer system.
【0070】
In the above explanation, the application program 22 is executed on the computer 211, but even if it is executed on the computers 212 and 213, the class is similarly distributed so that the load is automatically distributed according to the state of the computer system at that time. Is assigned and the same execution result is obtained.
【0071】
Further, in the above description, the application program written in the programming language C ++ has been described, but the application program written in another object-oriented language such as the programming language known as smalltalk has been described. The mechanism shown in this embodiment can also be applied.
【0072】
As described above, according to the present embodiment, according to the present invention, the load is automatically assigned to each computer constituting the distributed computer system for each method part of the class so as to distribute the load when the application program is executed. Therefore, it is not necessary to prepare the entire application or a part thereof in each computer in advance. In addition, all the methods assigned to another computer are deleted when the application program ends, so there is no waste of resources.
【0073】
[Effect of the invention]
As described above, according to the present invention, it is possible to provide a distributed computer system capable of sufficiently distributing the load of executing an application program without increasing the resource amount of each computer required for load distribution. Can be done.
[Simple explanation of drawings]
[Figure 1]
It is a block diagram which shows the structure of the distributed computer system which concerns on one Example of this invention.
[Figure 2]
It is explanatory drawing which shows the structure of an application program.
[Fig. 3]
It is explanatory drawing which shows the structure of a class.
[Fig. 4]
It is a block diagram which shows the structure of the load distribution means which concerns on one Example of this invention.
[Fig. 5]
It is explanatory drawing which shows the table stored in the load management part which concerns on one Example of this invention.
[Fig. 6]
It is explanatory drawing which shows the table stored in the class management part which concerns on one Example of this invention.
[Fig. 7]
It is a flowchart which shows the processing procedure performed by the execution selection means which concerns on one Example of this invention.
[Fig. 8]
It is a flowchart which shows the processing procedure of the class transfer means which concerns on one Example of this invention.
[Fig. 9]
It is a flowchart which shows the processing procedure performed by the class storage means which concerns on one Example of this invention.
[Fig. 10]
It is a flowchart which shows the processing procedure performed by the method execution means class storage means which concerns on one Example of this invention.
[Fig. 11]
It is a flowchart which shows the processing procedure performed by the method execution request means which concerns on one Example of this invention.
[Fig. 12]
It is a flowchart which shows the processing procedure performed by the method surrogate execution means which concerns on one Example of this invention.
[Fig. 13]
It is a flowchart which shows the processing procedure performed by the program termination means which concerns on one Example of this invention.
[Fig. 14]
It is a flowchart which shows the processing procedure performed by the class erasing means which concerns on one Example of this invention.
[Fig. 15]
It is a flowchart which shows the processing procedure performed by the class deletion request means which concerns on one Example of this invention.
[Fig. 16]
It is a flowchart which shows the application program execution procedure of the operation system which concerns on one Example of this invention.
[Fig. 17]
It is explanatory drawing which shows the state just before the application program execution of the distributed computer system which concerns on one Example of this invention.
[Fig. 18]
It is explanatory drawing which shows the state in the process of executing the application program of the distributed computer system which concerns on one Example of this invention.
[Fig. 19]
It is explanatory drawing which shows the state in the process of executing the application program of the distributed computer system which concerns on one Example of this invention. It is a state diagram at one time of the distributed computer system.
[Fig. 20]
It is explanatory drawing which shows the state in the process of executing the application program of the distributed computer system which concerns on one Example of this invention. It is a state diagram at one time of the distributed computer system.
[Fig. 21]
It is explanatory drawing which shows the state in the process of executing the application program of the distributed computer system which concerns on one Example of this invention. It is a state diagram at one time of the distributed computer system.
[Fig. 22]
It is explanatory drawing which shows the state in the process of executing the application program of the distributed computer system which concerns on one Example of this invention.
[Fig. 23]
It is explanatory drawing which shows the state after the application program execution completion of the distributed computer system which concerns on one Example of this invention.
[Explanation of symbols]
20 ...... Communication medium 22 ...... Application program 23 ...... Load distribution means 30 ...... Communication control means 31 ...... Memory 32 ...... Class Management Department 33 ...... Method execution method 34 ...... Method execution request means 35 ...... Class transfer means 36 ...... Class deletion request means 37 ...... Method proxy execution method 38 ...... Class storage means 39 ...... Class elimination means 40 ...... Load management department 41 ...... Execution selection means 42 ...... Program termination means 211, 212, 213 ...... Computer 231, 232, 233 ...... Load distribution means
Every citation, both ways
| Document | Relation | Office | Cited during |
|---|---|---|---|
| JPH08235127A | Cited by | Japan | Search report |
| US6421736B1 | Cited by | United States of America | Applicant |
| US6886167B1 | Cited by | United States of America | Applicant |
| US6345311B1 | Cited by | United States of America | Applicant |
| JPH09212364A | Cited by | Japan | Search report |
| KR100331519B1 | Cited by | Republic of Korea | Search report |
| US6895587B2 | Cited by | United States of America | Applicant |
3 priority claims, no other members on record
Priority claims3
| Document | Office | Kind | Date |
|---|---|---|---|
| 13604993 | Japan | A | |
| 5136049 | – | – | – |
| JP19930136049 | – | – | – |
Numbers
- Publication
- 6-348666
- Publication, DOCDB
- H06348666
- Publication, EPODOC
- JPH06348666
- Application
- 5136049
- Application, DOCDB
- 13604993
- Application, EPODOC
- JP19930136049
Titles3
- English
- PROGRAM EXECUTION LOAD DISTRIBUTION METHOD IN COMPUTER SYSTEM
- Japanese
- 【発明の名称】計算機システムにおけるプログラム実行負荷分散方法
- English
- [Title of Invention] A method for distributing program execution load in a computer system
Classification
- IPC, 3
- G06F15 16
- G06F9 44
- G06F15 177