System for managing execution of programs by multiple computing systems
Abstract
Problem to be solved.To quickly acquire a program to be executed. One or more of a program, to some extent, based on the location of one or more previously stored copies of the program from which a copy of the program to be executed can be obtained. Select the appropriate computer system on which to run the instance. For example, in some situations, the selection of the appropriate computer system to run an instance of a program is, such as a stored copy of the program, a running copy of the program, and / or an available computer system. Based to some extent on physical or logical proximity to other resources. [Selection diagram] Fig. 5

Term
5.5 yearsto projected expiry
Projected expiry 8 March 2032, counted from filing; an application has no term until it is granted.
- Priority
- Filed
- Published
- Today
- Projected expiry
1 claim: 1 independent, 0 dependent
- 1プログラムの実行を管理するように構成されたシステムであって、 メモリと、 指示されたプログラムの少なくとも1つのインスタンスを実行する決定に応答して、前記指示されたプログラムの1つまたは複数のインスタンスを実行するための1つまたは複数の計算機を自動的に選択するように構成されたシステムマネージャモジュールであって、前記選択された計算機は、複数のグループのうちの1つまたは複数のグループにあり、計算機を前記選択することは、前記選択された計算機が、前記指示されたプログラムの格納されたコピーを有する1つまたは複数のマシンを含むグループのメンバであることに少なくとも或る程度、基づく、システムマネージャモジュールと を含むことを特徴とするシステム。
57 paragraphs, as filed
The present invention generally relates to a program on a plurality of computer systems, such as performing a copy of a program between groups of computer systems to allow efficient acquisition of a program copy to be performed. Regarding managing execution.
A significant number of private data centers operated by a single organization and on behalf of a single organization, and public data centers that provide customers with access to computer resources under different business models. Data centers that house computer systems connected to each other have become commonplace. For example, some public data center operators provide network access facilities, power facilities and secure installation facilities for hardware owned by various customers, while other public data center operators. Provides a "full service" facility that also includes the actual hardware resources used by the customer. However, as the size and scope of traditional data centers grows, the task of providing, managing, and operating physical computer resources becomes increasingly complex.
The advent of virtualization technology for product hardware has provided a partial solution to the problem of managing large computer resources for many customers with diverse needs. Generally speaking, virtualization technology allows various computer resources to be shared efficiently and securely among multiple customers. For example, virtualization technologies such as those provided by VMWare, XEN, or User-Mode Linux are single physical computing. Machine) resources can be allowed to be shared among multiple users. More specifically, each user can be provided with one or more virtual machines hosted by a single physical computer, and each such virtual machine is separate. It is a software simulation that acts as a logical computer system. Each virtual machine provides the user with the experience of being the sole operator and administrator of a given hardware computer resource, while also providing application isolation and application security among the various virtual machines. To do. In addition, some virtualization technologies, such as a single virtual machine with multiple virtual processors that actually span multiple separate physical computer systems, provide virtual resources that span one or more physical resources. can do.
However, one problem that arises in the context of a data center that virtually or physically hosts a large number of applications or systems for a diverse set of users is managing the storage, distribution, and acquisition of copies of software applications. Is involved. Applications can, for example, be extremely large and have sufficient storage resources (if not impossible) to store local copies of all hosted applications on all computer systems in the data center. It can be expensive. However, if a centralized storage location where a copy of the application is frequently transmitted to all computer systems in the data center where the application should run is maintained as an alternative, it is still high in terms of network bandwidth resources. It is a cost. In such alternatives, network bandwidth is monopolized for application copy transmission and can prevent the running application from receiving sufficient network bandwidth for the application to operate. Further, there may be a considerable startup wait time for application execution, such as while waiting for application copy transmission to be reached. Such difficulties can be exacerbated by a variety of factors, such as the frequent introduction of new applications to be executed and / or the frequent deployment of successive versions of the application.
<p> Therefore, in order to address the above problems, it would be beneficial to efficiently distribute a copy of the application to the computer system running the application and to provide technology to provide various other benefits. ..</p>
<p> Techniques for managing program execution on multiple computer systems are described. In some embodiments, the techniques described are performed for a program execution service for executing multiple programs for multiple users (eg, customers) of the service. In some embodiments, the program execution service can use a variety of factors to select the appropriate computer system on which to execute an instance of the program. The factor is that the selected computer system can obtain a copy of the program to be executed and / or a copy of the available computer system resources for the execution of the program instance, before the program. Contains information such as the location of one or more stored copies. For example, in some embodiments, the choice of the appropriate computer system on which to run an instance of a program may be based to some extent on identifying the computer system that already contains a local copy of the program. is there. In another example, the choice of the appropriate computer system is one or more such local copies, such as one or more other computer systems in a common group with the identified computer system. It is somewhat possible to be based on identifying computer systems that are close enough (geographically and / or logically) to multiple computer systems.</p><p> In some embodiments, the plurality of computer systems available to execute a program are such as one or more networks capable of transmitting data between computers, or each other via another data exchange medium. It is possible to include multiple connected physical computers. Multiple computer systems can be located, for example, in a physical location (eg, a data center), can be divided into multiple groups, and multiple computer systems thereof. Can be managed by one or more system manager modules that are responsible for as a whole, and to manage the computer systems of a group by multiple machine manager modules, each associated with one of the groups. It is possible to do. At least some of the calculators each have enough resources to run multiple programs at the same time (eg, enough writable memory and / or enough storage, enough CPU cycles or other CPU usage indicators, enough network). It is possible to include one or more) such as bandwidth, sufficient swap space, etc. For example, at least some of the calculators in some such embodiments can each host multiple virtual machine nodes, each capable of running one or more programs for different users. .. As mentioned above, in at least some embodiments, the plurality of computer systems managed by the program execution service are based on physical proximity or logical proximity, criteria including having a common data exchange medium, and the like. , Can be organized into multiple separate groups (eg, each computer system belongs to a single group). In one example, a common data exchange medium for a group can be provided by a single network switch and / or rackbackplane that provides high bandwidth communication between the group's computer systems (eg). For example, some or all of the computer systems connected to the network switch or rack backplane are members of the group). Also, each group of computer systems is based on one or more other data exchange media, such as other data exchange media that have a lower bandwidth than the group's common data exchange media (eg, Ethernet®). It is also possible to connect to another computer system (eg, another group of computer systems, or a remote computer system that is not managed by a program execution service) by wiring, wireless connections, or other data connections. Further, in at least some embodiments, some or all of the computer systems can be used to store a local copy of the executable program, such as before or during program execution, respectively (eg, a local program repository). , Hard disk, or other local storage mechanism), respectively. Further, in at least some embodiments, each group of multiple computer systems uses one or more computer systems in that group to be used locally by the other computer systems in that group. A copy can be stored. A local program repository (eg, a hard disk, or other local storage mechanism) that some or all of the data systems can use to store a local copy of the program, such as before or during program execution, respectively. It is possible to have each. Further, in at least some embodiments, each group of multiple computer systems uses one or more computer systems in that group to be used locally by the other computer systems in that group. A copy can be stored. A local program repository (eg, a hard disk, or other local storage mechanism) that some or all of the data systems can use to store a local copy of the program, such as before or during program execution, respectively. It is possible to have each. Further, in at least some embodiments, each group of multiple computer systems uses one or more computer systems in that group to be used locally by the other computer systems in that group. A copy can be stored.</p><p> In an exemplary embodiment, the program execution service can include a software facility that runs on one or more computer systems to manage the execution of the program. A software facility may include, for each group of one or more computer systems, one or more machine manager modules that manage the acquisition, storage, and execution of programs by that group of computer systems. For example, each separate computer can be provided with a separate machine manager module, such as a machine manager module for a computer, running on at least one virtual machine in multiple virtual machines on that calculator. Is. Further, in some embodiments, the software facility can include one or more system manager modules running on one or more computer systems, which system manager modules execute the program. Manage the acquisition, storage and execution of programs for all of the multiple computer systems used to do so. The system manager module can interact with the machine manager module as appropriate, as described below.</p><p> In at least some embodiments, execution of one or more instances of a program on one or more computer systems can be initiated in response to a current execution request for immediate execution of those program instances. Is. Alternatively, this execution start may be based on a previously received program execution request that is scheduled or reserved for execution of a future program instance at this time. Program execution requests can be received in various ways. For example, it may be directly from the user (eg, through an interactive console or through another GUI provided by a program execution service), or another program, or one or more of the program itself. It may be from a running program of the user that automatically starts the execution of the instance (via an API (Application Program Interface) provided by the program execution service, for example, an API that uses a web service).</p><p> A program execution request can contain various information that should be used when starting execution of one or more instances of a program, such information, for example, for future execution. A single desired instance of an instance, such as the number of previously registered or otherwise supplied program instructions and the number of instances of the program to be executed concurrently (eg, the minimum and maximum number of desired instances, etc.) (Represented as a number) and so on. In addition, in some embodiments, the program execution request can include various other types of information. For example, the information may allow the user account instructions, or other instructions of a previously registered user (eg, when identifying a previously stored program, to execute additional / or requested program instances. Instructions to use when determining if there is), payment source instructions to use when providing payment for the program instance execution to the program execution service, previous payment authorization for the program instance execution, or other Permission instructions (eg, previously purchased contracts that are valid for a certain amount of resource utilization for a certain number of program execution instances over a period of time), and / or are executed immediately. There should be an executable copy of the program, or any other copy that should be stored for further / or later execution. Further, in some embodiments, the program execution request may further include various other types of preferences and / or requirements for the execution of one or more program instances. Such preferences and / or requirements are indicated in one of a plurality of data centers accommodating multiple available computers, on multiple computer systems close to each other, and / or in one or more other. A computer system that runs a program instance (for example, an instance of the same program or an instance of another program) It is possible to include instructions that some or all of the program instances should run at the indicated geographic and / or logical location, such as on one or more computer systems close to the program. is there. Such preferences and / or requirements can further include instructions that some or all of the program instances should be assigned the indicated resources, respectively, during execution.</p><p> After receiving a request to run one or more instances of a program at the point in time, the program execution service determines which computer system to use to run those program instances. In some embodiments, the determination of the computer system to be used is made at the time of the request, even if it is for future execution. In other embodiments, the determination of the computer system to be used for future execution of one or more program instances is at a later point in time, such as at future execution, based on the information available at that time. It is possible to postpone until. The determination of which computer system should be used to run each program instance is specified elsewhere in the program request, or for the program and / or associated users (eg, for example). It can be done in a variety of ways, including based on any preference and / or requirement (at the time of previous registration). For example, if criteria for preferred and / or required resources for running a program instance have been determined, then the computer system determines the appropriate computer system to run the program instance with those resource criteria. It can be at least to some extent based on the availability of sufficient resources to meet.</p><p> In some embodiments, the program execution service determines which computer system should be used to execute the program to be executed, one stored before the program to be executed or It can be based on the location of multiple copies. In particular, as mentioned above, in at least some embodiments, the various computer systems available to execute the program can be organized into groups (each computer system can be in multiple groups). Belonging to one, etc.). Therefore, the determination of whether a computer system is appropriate to run an instance of a program is whether one or more computer systems in the group of computer systems store a local copy of the program. Somehow at least it is possible to be based. To run an instance of a program, the program by selecting a computer system that already has a locally stored copy of the program, or a computer system that belongs to a group that has one or more locally stored copies. Various benefits can be obtained, such as shortening the waiting time for program execution based on obtaining a copy of. When a computer system in a group stores a local copy of a program to be executed, a program execution service can be used, for example, for a computer system with a locally stored copy to execute an instance of the program. An instance of a program for a variety of reasons, such as when you do not currently have sufficient resources, or when a computer system with a locally stored copy is already running one or more instances of the program. It is possible to select one or more other computer systems in the group to currently run.</p><p> In a further embodiment, the program execution service may choose one or more computer systems to execute an instance of the program due to various other factors. For example, if a user requests that multiple instances of the indicated program be executed at the same time, the program execution service distributes the execution of these program instances among computer systems that are members of different groups. You can choose to bring higher reliability in the face of group-specific network outages, or other problems. Similarly, in some embodiments, it is possible for multiple instances of a program to run on multiple computer systems rather than on a single computer system (that single computer system is them). Even if you have enough resources to run multiple instances of. Higher in the face of such distribution of program instances, for example, the failure of a single computer system running all of the program instances instead, or the loss of connectivity to that single computer system. It is possible to bring reliability. In addition, if the computer systems managed by the program execution service are physically (eg, geographically) separated, then the program execution service is a program on a computer system located within a single data center. It can be instructed by the user to run multiple instances to provide relatively high network bandwidth for communication between instances of the running program, or otherwise. It is possible to choose. Alternatively, the program execution service may provide several separate end users or several separate program instances that are geographically dispersed, if the program instances have little or no interaction with each other, and / or various program instances. Multiple, such as when supporting multiple applications</p><p> After deciding which computer system to use when running an instance of a program, the program execution service can initiate the execution of those program instances in a variety of ways. For example, the system manager module can give instructions and various other execution information to the selected computer system. Other such information may include, for example, instructions of one or more other computer systems that store or may store a local copy of the program. Other types of information given to the selected computer system include instructions on how long to run the program instance, instructions on the resources that should be assigned to the program instance, and instructions on the permissions that should be given to the program instance. , Can include restriction instructions on how the execution of the program instance should be managed (eg, what type of communication should the program instance be allowed to send or receive). is there.</p><p> After the selected computer system is notified to execute one or more instances of the indicated program, the selected computer system receives instructions or other relevant information (eg, predefined). Attempts to perform program instance execution according to the preferred preferences or requirements). In at least some embodiments, this program execution notification belongs to a machine manager module associated with the selected computer system (eg, a machine manager module running on the selected computer system, or a selected computer system). It can be received by the machine manager module) that runs for the group. In such an embodiment, the machine manager module can act to manage the execution of program instances. For example, in situations where the selected computer system does not already store a local copy of the indicated program to be executed, the machine manager module will perform it for execution as well as for storing options locally. It is possible to get a copy of the program or act to get it elsewhere. One known to obtain a program copy, for example, is indicated in a notification or may have stored at least a local copy of the program to request or obtain a copy of the program. Or it can include contacting multiple computer systems or other systems (eg, data storage systems). Acquisition of a program copy can be performed in different ways in different embodiments, including by receiving a copy of the program along with a notification received to execute the program instance, as described below. .. As described below, the program execution service performs various other actions that control the execution of the program, at least in some embodiments.</p><p> In another aspect, another program may programmatically initiate a request to execute a program instance and, optionally, perform various other types of management, provision, and operational actions. APIs that enable it can be provided. Such operations include creating user accounts, reserving execution resources, registering new programs to be executed, managing groups and access policies, monitoring and managing running program instances, and so on. Not limited to. The functionality provided by the API can be invoked by the client computer system and the client device for the user, including, for example, by a program instance running on the computer system of the program execution service.</p><p> Some embodiments in which the execution of a particular type of program on a particular type of computer system is managed in a particular way will be described later for illustration. These examples are provided for illustration purposes and are simplified for simplicity. The techniques of the present invention can also be used in a wide variety of other situations, some of which are described later, and these techniques can be used in virtual machines, data centers, or other specific types of computer systems or computer systems. Not limited to being used in the configuration.</p>
<figref num="1">FIG. 5 is a network diagram illustrating an exemplary embodiment in which a plurality of computer systems exchange and execute programs.</figref><figref num="2">FIG. 5 shows an example of a group of computer systems that stores and exchanges copies of a program.</figref><figref num="3">It is a block diagram which shows an exemplary computer system suitable for managing the execution of a program on a plurality of computer systems.</figref><figref num="4A">It is a flow chart which shows the embodiment of the system manager module routine.</figref><figref num="4B">It is a flow chart which shows the embodiment of the system manager module routine.</figref><figref num="5">It is a flow chart which shows the embodiment of the machine manager module routine.</figref><figref num="6">It is a flow chart which shows the embodiment of the program execution service client routine.</figref>
FIG. 1 is a network diagram showing an exemplary embodiment in which a plurality of computer systems exchange and execute programs under the control of a program execution service or the like. In particular, in this example, the program execution service manages the execution of programs on various computer systems located within the data center 100. The data center 100 includes a plurality of racks 105, each rack comprising a plurality of computer systems 110a-c, and in this exemplary embodiment also including a rack support computer system 122. Each of the computer systems 110a-c hosts one or more virtual machines 120, as well as a separate node manager 115 that manages these virtual machines. Each virtual machine 120 can be used to provide an independent computer environment for running instances of the program. The rack support computer system 122 can provide various utility services to racks and, optionally, to other computer systems that are local to other computer systems located within the data center. Utility services can include, for example, data storage and / or program storage for other computer systems, execution of one or more machine manager modules that support other computer systems, and so on. Each computer system 110 can, as an alternative, have a separate machine manager module (eg, provided as part of the computer system's node manager) and / or local storage to store a local copy of the program. It is possible to have (not shown). The computer systems 110a-c and the rack-supported computer systems 122 all share a common data exchange medium and can all be part of a single group. This common data exchange medium can be used, for example, in another data center 100.
Further, the exemplary data center 100 further includes computer systems 130a-b and 135 that share a common data exchange medium with node manager 125, which manages computer systems 130a-b and 135. In addition, the computer system 135 hosts a plurality of virtual machines as an execution environment to be used when executing a program instance for one or more users, while the computer systems 130a to 130b are separate. Do not host virtual machines. An optional computer system 145 is present at the interconnection point between the data center 100 and the external network 170. The optional computer system 145 can provide several services, such as acting as a network proxy and managing incoming and / or outgoing data transmissions. In addition, an option to help manage the execution of programs on other computer systems located within the data center (or optionally on computer systems located in one or more other data centers 160). System Manager Computer System 140 is also illustrated. Optional System Manager Computer System 140 can run the System Manager module. As mentioned earlier, in addition to managing program execution, the System Manager module manages user accounts (eg, create, delete, billing, etc.), register, store, and distribute programs to be executed, programs. A variety of services can be provided, including collecting and processing performance and audit data related to the execution of a program, and getting paid for the execution of a program from a customer or other user.
The data center 100 includes a computer system 180 and a data center 160 that can be operated by the operator or a third party of the data center 100 via a network 170 (eg, the Internet), as well as an optional system manager 150. , Connected to several other systems. In a manner similar to System Manager 140, System Manager 150, in addition to providing various other services, executes programs on computer systems located in one or more data centers 100 and / or 160. Can be managed. Although the exemplary system manager 150 is shown as being external to any particular data center, in other embodiments the system manager 150 is within a data center, such as one of the data centers 160. May be placed in.
Figure 2 shows an example of two groups of computer systems that store and exchange copies of a program, such as for program execution services. It will be appreciated that in actual embodiments, the number of groups, computer systems, and programs may be much larger than the groups shown in FIG. For example, in one exemplary embodiment, there could be 40 computer systems per group, 100 groups per data center, and 4000 computer systems per data center. Yes, each computer system can host 15 virtual machines running a customer's program instance. In addition, if each group contains a dedicated computer system with 2 terabytes (TB) of storage, 2000 1GB virtual machine image program copies can be stored per group, in which case total per data center. It will be 200,000 copies. Alternatively, if each of the 40 computer systems per group has 100GB of local storage, then 4000 1GB virtual machine image program copies can be stored per group, for a total of 400,000 per data center. It will be a copy of. If each hosted virtual machine runs one program, such a data center can run as many as 6000 program instances at the same time. In other embodiments, other numbers of groups, other numbers of computer systems, and other numbers of programs can be used, and much smaller and / or variable size programs. Please understand that it can be stored and executed.
FIG. 2 shows two groups, group A200 and group B250. Group A includes three computers 210a-c named MA1, MA2 and MA3 respectively. Group B250 also includes three computers 260a-c named MB1, MB2 and MB3 respectively. Each group may have a different number of calculators of different types, and in some embodiments the calculators may be members of multiple groups or not members of any group. .. As detailed elsewhere, the computers in each group share a common data exchange medium for that group (not shown).
In an exemplary example, each computer in Figure 2 has one or more local program copies in a local program repository (eg, a repository as part of persistent storage provided by a hard disk or other storage device). Can be stored. Computer MA1 has a local copy of programs P1, P2, P3, P5 and P9 stored in the program repository 220a of computer MA1 and currently runs an instance of program P1 as shown in box 230a. doing. The storage capacity of the program repository on each computer is limited to a maximum of five program copies, and the execution resources of each computer system are limited to a maximum of two program instances that are executed simultaneously. The size of the program repository used in this example, and the limit on the number of programs to be executed, is merely an example, and in other embodiments, each computer system may have additional separate resources. In addition, the size of the program repository can, in many embodiments, be one or several orders of magnitude larger than the size of memory available to run the program instance. That is not always the case. In other embodiments, the maximum number of concurrently running program instances may be greater than, less than, or less than the number of program copies that can be stored locally in the program repository, or with such a number. It may be the same. Therefore, at least some computers or other systems may provide only one of the local program repositories and the available resources to run the program instances. Finally, as detailed elsewhere, in some embodiments, a locally stored copy of at least some of the programs will be made of other program copies after the program repository has reached capacity. Some to make room for Under circumstances, they may be evacuated from storage or excluded. In some embodiments, at least some of the programs' running instances have space for other running program instances after the program execution resource has reached its capacity. Under some circumstances, the run may be terminated or excluded from the run.
A plurality of exemplary scenarios are presented herein for exemplary purposes to provide examples of some types of behavior in one embodiment of a program execution service. Program execution services can influence the placement of running program instances on a computer using one or more specified, predefined, and / or learned policies, such as: , A simplified set of policies is used in this example. First, multiple instances of the program run on multiple groups of computers, if possible. Second, multiple instances of the program run on multiple computers, if possible. Third, an instance of a program runs on a computer that already has a copy of the program in the program repository, if possible. Fourth, an instance of a program runs on a computer that is a member of a group that has at least one computer that stores a local copy of the program in the program repository, if possible. Finally, the instance of the program runs on the computer with the most available execution resources, if possible.
From the exemplary example of managing program execution for these six computer systems, assume that the client of the program execution service requested the execution of two instances of program P7. In this case, given the policy described above, an exemplary embodiment of the program execution service is to execute one instance of P7 in group A and one instance of P7 in group B. Is likely to be selected. This is because such arrangements tend to distribute copies across multiple groups. Between the computers in group A, none of the computers in this group store a local copy of this program, so the program execution service is because the computer MA3 is already running two programs (P8 and P9). It is likely that you will not choose to perform a copy of P7 on computer MA3. Between computer MA1 and computer MA2, MA2 is selected for execution because MA2 is not currently executing any program. In one embodiment, computer MA2 obtains a copy of program P7 from one or more computer systems outside Group A for execution and, optionally, for storage locally in repository 220b. .. For example, the computer MA2 can obtain a copy of program P7 from a remote program repository for all of the programs in the program execution service and / or from a location external to the program execution service. For Group B computers, the program execution service can choose any of the three computers to run the P7 program instance. This is because none of these computer systems store a local copy of this program, and each of these computers is running one program. However, the program execution service is currently 1 in the program repository of computer MB3. Since it stores only one program copy, the computer MB3 can be selected. Therefore, the computer MB3 can store a local copy of the program P7, if desired, without having to evict the program copy stored from the program repository of the computer MB3.
Next, again, suppose that the client of the program execution service requested the execution of two instances of program P6, starting from the initial conditions shown in Fig. 2. In this case, given the policy described above, an exemplary embodiment of the program execution service is to execute one instance of P6 in group A and one instance of P6 in group B. Is likely to be selected. This is because such an arrangement distributes the instances across multiple groups. Among Group A computers, calculator MA2 may still be selected, but none of these computer systems stores a local copy of program P6, and calculator MA2 is the least busy. Therefore, it is expensive. Among the equally busy computers in Group B, only MB2 stores a local copy of that program because of the policy of preferring to distribute a copy of a single program across multiple computers in the group. Despite the fact that it is, the computer MB2 may not be selected. However, another embodiment with a different policy that reflects more efficiency than reliability is to run P6 on the computer MB2 because a copy of P6 is already stored in the MB2 program repository. Note that it is possible to actually select. Between the remaining candidate computers MB3 and MB1, the program execution service can prefer the computer MB3 because there is no need to evict a copy of the program from the MB3 program repository as a possibility. Therefore, MB3 obtains a copy of program P6 from MB2 for execution and optionally storage in the local repository 270c in this embodiment.
Next, again, suppose that the client of the program execution service requested the execution of one instance of program P4, starting from the initial conditions shown in Fig. 2. In this case, given the policy described above, an exemplary embodiment of the program execution service is likely to choose to execute P4 on computer MB1. In particular, there are no instances of P4 already running, and only one instance was required to run, so a policy that favors distributing program instances among multiple groups. , And policies that prefer to avoid placing multiple running instances of a program on a single computer do not apply. Therefore, MB1 already stores a local copy of program P4 in MB1's program repository, so MB1 is likely to be chosen to run P4.
Next, again, suppose that the client of the program execution service requested the execution of one instance of program P10, starting from the initial conditions shown in Fig. 2. In this case, given the policy described above, an exemplary embodiment of the program execution service is likely to choose to run P10 on MA2. As in the previous example, policies that prefer to distribute instances of the program for execution among multiple groups, and avoid having multiple instances of the program on a single computer. The policy does not apply. Computer MA3 is an attractive candidate because it already has a copy of P10 in the computer MA3 repository, but it has already reached the limit of computer MA3, two programs to be executed (P8 and P9). Therefore, it does not currently have the capacity to run P10. As a result, the computers MA1 and MA2 are in the same group as the computer (MA3) that has a local copy of the program P10 stored in the repository, which is preferable to any of the computers in group B. Be left behind. Between MA1 and MA2, MA2 is the least likely to be selected because it is the least busy, and MA2 gets a copy of program P10 from MA3.
Then again, starting with the initial conditions shown in Figure 2, assume that the client of the exemplary embodiment of the program execution service requested the execution of six more instances of program P3. In this case, given the policy described above, the program execution service is likely to run two instances on computer MA2 and one instance on each of the computers MA1, MB1, MB2 and MB3. Since the computer MA3 has already reached the limit of the computer MA3, which is two programs to be executed (P8 and P9), it is highly likely that the instance will not be executed on this computer at all. In this case, in some embodiments, it is possible to forcibly evict the stored local copy of the program from a computer having a program repository that does not have extra capacity in order to store the local copy of the program P3. Please note that there is. For example, in an embodiment in which a copy of a program to be executed is always selected to be stored in a local program repository prior to execution, computers MA1 and MB1 force one local program copy from their respective program repositories. It is possible to move out. It should also be noted that in this case, the computers MA2 and MB3 will each execute two instances of P3, contrary to the policy of preferring to distribute multiple instances of the program to be executed across multiple computers. I want to. However, since there is no additional computer to run the P3 program instance in this given example, the program execution service will run multiple instances of P3 on a single computer if that requirement is met. You will choose to do it. Alternatively, in some embodiments, the program execution service applies different weights to the policy so that the program execution service runs a single instance on each of the computers MA1, MA2, MB1 and MB3. , Less than the requested number of instances It is possible to choose to do it instead. Similarly, in some embodiments, if more than six additional instances of program P3 are requested and the program and / or requester has a sufficiently high priority, the program execution service will be another program instance ( For example, terminating the execution of programs P8 and / or P9 instances on MA3, and / or after one of the currently running program instances spontaneously terminates, is then available for P3. You may choose to run more instances of P3 instead, such as by reserving program instance execution.
Continuing with the current example, computer MB1 stores a local copy of the program together with MB2 and MB3 from group B, as do computers MA1 and MA2 in group A, so for execution of program P3 Has multiple available resources to obtain a copy. In this embodiment, MB1 is the first X bytes and the second X bytes, where both MB2 and MB3 in MB1's own group are part of program P3 (eg, X is the number selected by the program execution service). ) Is provided. The calculator MB1 then monitors how quickly the response is received from those calculators and is more responsive to supply at least most (and possibly all) of the rest of the program. Require a good calculator. In other embodiments, obtaining a copy of program P3 for computer MB1 requires a program copy from either computer MB2 or MB3, computers MA1 and / or MA2 in group A (MB2 in group B). And may be executed in other ways, such as by requesting at least some part of the program copy from (in addition to MB3, or instead of MB2 and MB3).
FIG. 3 is a block diagram showing an exemplary computer system suitable for managing the execution of a program on a plurality of managed computer systems, such as by executing an embodiment of a program execution service system. is there. In this example, the computer system 300 executes certain embodiments of the system manager module to coordinate the execution of programs on a plurality of managed computer systems. In some embodiments, the computer system 300 can correspond to the system manager 140 or 150 of FIG. In addition, one or more machine manager computer systems 370 each run machine manager module 382 to facilitate program acquisition and execution by one or more associated computer systems. In some embodiments, each of the one or more machine manager modules 382 can correspond to either node manager 115 or 125 of FIG. In this example, multiple machine manager computer systems 370 are provided, and each system 370 operates as one of multiple computer systems of the program execution service managed by the system manager module. In the illustrated example, a separate machine manager module 382 is run on each system of computer system 370. In other embodiments, the machine manager module 382 on each system of the machine manager computer system 370 is capable of managing one or more other computer systems (eg, another computer system 388) instead. ..
In this exemplary embodiment, the computer system 300 includes a "CPU" (central processing unit) 335, a storage 340, a memory 345, and various "I / O" (input / output) devices 305. The illustrated I / O devices include a display 310, a network connection 315, a computer readable media drive 320, and other I / O devices 330. Other I / O devices not shown may include keyboards, mice or other pointing devices, microphones, speakers, and the like. In the illustrated embodiment, the system manager module 350 is running in memory 345 to manage the execution of programs on other computer systems, with one or more other programs 355 as options. , May be executed in memory 345. The computer system 300 and the computer system 370 are connected to each other and also to other computer systems 388 via network 386.
Each computer system 370 also includes a CPU 374, various I / O devices 372, a storage 376, and a memory 380. In the illustrated embodiment, the machine manager module 382 is a memory for managing the execution of one or more other programs 384 on a computer system for a program execution service, such as for a customer of the program execution service. Running in 380. In some embodiments, some or all of the computer system 370 may host multiple virtual machines. If so, each of the running programs 384 is an entire virtual machine image (eg, having an operating system and one or more application programs) running on a separate hosted virtual machine. May be good. The machine manager module 382 can run similarly on another hosted virtual machine, such as a privileged virtual machine that can monitor other hosted virtual machines. In other embodiments, the running program instance 384 and the machine manager module 382 can run as separate processes on a single operating system (not shown) running on computer system 370. Thus, in this exemplary embodiment, the capabilities of the program execution service communicate over network 386 to jointly manage the distribution, acquisition, and execution of programs on the managed computer system System Manager 350. And provided by an interaction with the machine manager module 382.
It will be appreciated that computer systems such as computer systems 300 and 370 are merely exemplary and are not intended to limit the scope of the invention. Computer systems 300 and 370 can connect to other devices not shown, including network-accessible database systems or other data storage devices. More generally, computers or computer systems or data storage systems are, without limitation, desktop or laptop computers, database servers, network storage devices and other network devices, PDA, cell phones, wireless phones, pagers, electronic organizers. , Internet equipment, TV-based systems (eg, using set-top boxes and / or personal / digital video recorders), and various other consumer products, including appropriate intercommunication capabilities, interact and explain It may include any combination of hardware or software capable of performing any type of function. In addition, the functionality provided by the illustrated system modules may be combined into fewer modules or distributed to more modules in some embodiments. Similarly, in some embodiments, the functionality of some of the illustrated modules may not be provided and / or other additional functionality may be available.
Also, various items are exemplified as being stored in memory or in storage while in use, but these items, and parts of these items, are of memory management and data integrity. It will also be appreciated that for purposes it is possible to transfer between memory and other storage devices. Alternatively, in other embodiments, some or all of the software components and / or software modules run in memory on another device and are exemplified via computer-to-computer communication. May communicate with. Some or all of the system modules or data structures can also be stored on computer-readable media such as hard disks, memory, networks, suitable drives, or portable media products that should be read via the appropriate connections. (For example, as a software instruction or structured data). Also, system modules and data structures can be generated as data signals (eg, carrier waves, or other analog or digital) on a variety of computer-readable transmission media, including wireless-based and wired / cable-based media. It can also be transmitted (as part of the transmitted signal) and in various forms (eg, as a single analog signal, or as part of a multiplexed analog signal, or multiple discrete digital packets. Or as a digital frame). Also, such computer program products can take other forms of other embodiments. Therefore, the present invention may be implemented in other computer system configurations.
4A to 4B show a flow chart of an embodiment of the system manager module routine 400. This routine is performed by the execution of the system manager module 140 and / or the system manager module 350 of FIG. 3, for example to manage the execution of multiple programs on multiple computer systems for program execution services. It can be provided.
This routine begins at step 405 and receives a status message or request related to the execution of one or more programs. The routine then proceeds to step 410 to identify the type of message or request received. If it is determined that a request to execute one or more instances of the indicated program is received, the routine proceeds to step 415. At step 415, the routine identifies one or more groups of computer systems to execute the indicated program. At step 420, the routine selects one or more computer systems within each group of one or more identified groups to execute an instance of the indicated program. This selection of one or more groups is whether the group has one or more computer systems that store one or more local copies of the program, the availability of appropriate computer resources, and them. It can be based on a variety of factors, such as the location of a group of computer systems. The choice of one or more computer systems within the identified group also depends on various factors, such as the location of the stored local copy of the program among the computer systems of that group, and the availability of computer resources. May be based on. As mentioned above, various specified policies and other criteria, including criteria specified by the user or other requester, can be used as part of the selection of groups and computer systems in various embodiments. Moreover, in other embodiments, the group and the particular computer system are individually selected, regardless of the group (eg, if the group is not used at all), such as simply selecting one or more of the most suitable computer systems. It is possible that it is not selected for.
Then, in step 425, the routine is selected computer systems, and / or those computer systems, such as by sending a message containing instructions for the programs to be executed, instructions to execute those program instances, and so on. Give to the machine manager module associated with. In an exemplary embodiment, a separate machine manager module runs on each system of the computer system and receives the message. As mentioned earlier, various types of information can be given to the machine manager module, including instructions on how to identify one or more computer systems for which a copy of the program to be executed should be obtained. It is possible. Alternatively, in some embodiments, the system manager supplies a copy of the indicated program directly to the computer system and / or programs on the computer system without the intervention of the machine manager module or other additional modules. It is possible to start the execution of.
If it is determined in step 410 that a request to register a new program has been received instead, such as from a user, the routine proceeds to step 440, instructing the program, as well as the identity of the user who registered the program, etc. Stores any relevant management information for. Then, at step 445, the routine optionally begins delivering a copy of the indicated program on one or more computer systems. For example, in some embodiments, the system manager may have one or more in one or more data centers having a stored copy of the indicated program in order to improve the efficiency of later program execution initiation. You can choose to seed your computer system and / or one or more repositories.
If, in step 410, it is determined that a status message has been received instead that reflects the behavior of one or more of the managed computer systems, the routine proceeds to step 450 and that one or more computers. Update status information about the system. For example, the machine manager module can determine that the associated computer system has modified the running program instance and / or the stored local program copy, and accordingly sends a status message to the system manager. Can be given. In some embodiments, the status message is to keep the system manager informed about the operational status of the managed computer system for use in selecting the appropriate computer system to run the program on. , Sent regularly by the machine manager module. In other embodiments, the status message can be sent at other times (eg, whenever relevant changes occur). In other embodiments, the system manager module can instead request information from the machine manager module, if desired. Status messages include the number and ID of programs currently running on a particular computer system, the number of copies of the program currently stored in the local program repository on a particular computer system, and ID, performance-related information about a computer system, and resource-related information (eg, utilization of CPU, network, disk, memory, etc.), configuration information about a computer system, and on a particular computer system. It can contain various types of information, such as reports of errors or failures related to hardware or software.
If, in step 410, it is instead determined that some other type of request has been received, the routine proceeds to step 455 and performs other indicated actions as appropriate. Such behavior includes, for example, responding to status queries from other components in the system, suspending or terminating the execution of one or more currently running programs, currently It can include moving a running program from one computer system to another, shutting down or restarting the system manager.
After steps 425, 445, 450 and 455, the routine proceeds to step 430 to calculate billing information about the user, update the display information, and send periodic queries to the node manager or other components. Perform any housekeeping task as an option, such as alternating logs or other information. The routine then proceeds to step 495 to determine whether to continue. If it continues, the routine returns to step 405, otherwise it proceeds to step 499 and returns.
FIG. 5 shows a flow diagram of the machine manager module routine 500. Routines are used, for example, to facilitate the acquisition of program copies for one or more related computer systems being managed, and the execution of program instances, for example, the machine manager module 382 in Figure 3, and / or It can be provided by running Node Manager 115 or 125 in Figure 1. In the illustrated embodiment, for a single computer system, each machine manager module routine is configured to run one or more program instances and also store one or more local program copies. The machine manager module works in concert with the system manager module routines described in relation to FIGS. 4A-4B to execute programs related to the managed computer system for program execution services. to manage.
The routine begins at step 505 and receives requests related to the execution of one or more programs, such as from the system manager module. The routine proceeds to step 510 to determine if a request to execute or store the indicated program has been received. If received, the routine proceeds to step 515 to determine if the indicated program is currently stored in the local program repository of the managed computer system. If not, the routine proceeds to step 540 to determine if the local program repository has enough capacity to store the indicated program. If there is not enough capacity, the routine proceeds to step 545 and is local as indicated in the request received in step 505, or otherwise based on the evacuation policy used by the machine manager module. Force one or more programs out of the program repository. If, after step 545, or at step 540, instead the local program repository is determined to have sufficient capacity to store a local copy of the indicated program, the routine proceeds to step 550 to identify the other. Obtain a copy of the indicated program from one or more computer systems. The routine can identify other computer systems that have a stored local copy of the program in various ways, including based on the information received as part of the request received in step 505. In addition, it is possible to use one or more other techniques such as broadcasting to neighboring computer systems, requesting central directories, and / or peer-to-peer data exchange. In another embodiment, a copy of the program can be provided with the request in step 505 instead. The routine then proceeds to step 555 to localize the acquired copy of the indicated program. Store in the gram repository. If, after step 555, or at step 515, it is determined that the indicated program was already in the repository, the routine proceeds to step 520 to determine if instructions for the program to be executed have been received. judge. If received, the routine proceeds to step 525 and begins executing the instructed program.
If, in step 510, it is determined that the request to store or execute the program was not received instead, the routine proceeds to step 535 and performs other indicated actions as appropriate. Other behaviors were collected regarding the performance of the program, such as responding to received requests and / or the program was running irregularly or over-utilizing resources. It may include suspending or terminating the execution of one or more programs, such as informed. In addition, other actions may include responding to requests for status information about the program currently running, or the contents of a local program repository.
If, instead, instructions for the program to be executed are not received at step 535, after step 525, or at step 520, the routine proceeds to step 530 and sends a status information message to one or more system manager modules. To send. In an exemplary embodiment, the routine sends a status information message to the system manager module each time an operation is performed to keep the system manager informed of the state of the computer system managed by the node manager. In other embodiments, the status information can be transmitted at other times and in other ways. After step 530, the routine proceeds to step 595 to determine whether to continue. If it continues, the routine returns to step 505, while if it does not, it proceeds to step 599 and returns. Although not shown herein, routines can also perform different housekeeping operations at different times, if desired.
FIG. 6 shows a flow chart of an embodiment of the program execution service client routine. Routines are provided by an application residing on one of the computer systems 180 shown in FIG. 1, such as providing an interactive console that allows a human user to interact with a program execution service. It is possible. The routine may, as an alternative, reflect the ability of the program execution service to be provided interactively to the user and / or programmatically to the user's program. Alternatively, this routine is part of one of the programs running by the program execution service on one of the managed computer systems, such programs being load balanced, increased or decreased. It is possible to dynamically execute additional program instances for the purpose of meeting demand.
The routine begins at step 605 and receives requests related to the execution of one or more programs. At step 610, the routine identifies the type of message received. If this request relates to the registration of a new program (or a new version of a previously registered program), the routine should proceed to step 625 and be registered with the program execution service (eg, the system manager module). Send instructions for new programs. This instruction can include a copy of the program, or instructions on how to acquire the program. If the request is determined in step 610 to instead be related to the execution of the program, the routine proceeds to step 615 and sends the request to the program execution service to execute one or more instances of the program to be executed. Send (for example, to the system manager module). For example, a routine can use instructions previously received from a program execution service to identify the program and / or user who wants the program instance to run. If, in step 610, instead it is determined that some other type of request has been received, the routine proceeds to step 625 and performs other indicated actions as appropriate. For example, a routine sends a request to the program execution service to reserve a computer resource to execute one or more indicated program instances at some future point in time, and the current or previous execution of one or more programs. Sends a status query about to the program execution service, provides and modifies user-related information (for example, as part of registering the user with the program execution service), unregisters or unregisters a previously registered program, You can suspend or terminate the execution of one or more program instances.
After step 615, 625 or 630, the routine proceeds to step 620 to optionally update the display information or in response to step 615, 625 or 630 with the information returned by the program execution service (not shown). Perform additional housekeeping tasks such as storing and periodically querying the status of the program execution service. After step 620, the routine proceeds to step 695 to determine whether to continue processing. The routine returns to step 605 if it continues, goes to step 699 if it does not continue, and returns.
Also, in some embodiments, the functionality provided by the routines described above can be provided in alternative ways, such as being divided into more routines or integrated into fewer routines. Those skilled in the art will understand. Similarly, in some embodiments, the illustrated routine may provide more or less functionality than described, while other illustrated routines may instead. It is possible if such functionality is lacking or included, or if the amount of functionality provided is changed. Further, various actions can be exemplified to be performed in a particular way (eg, sequentially or in parallel) and / or in a particular order, but others. It will be appreciated by those skilled in the art that in embodiments, the operations can be performed in other order and in other ways. In addition, the above-mentioned data structures can be structured in various ways, such as by dividing a single data structure into a plurality of data structures or by integrating a plurality of data structures into a single data structure. It will also be understood by those skilled in the art that it is possible. Similarly, in some embodiments, the illustrated data structure may store more or less information than described, or if other illustrated data structures lack such information or It is possible if it is included or if the amount or type of information stored is changed.
As mentioned above, various embodiments organize the computer systems of the program execution service into one or more groups to facilitate the implementation of the policies associated with the execution of the program. In addition, computer systems can be organized in other ways, such as by using a hierarchy of groups. For example, each of the smallest groups can contain a single computer system, and each computer system is assigned to its own group. In that case, a single machine group connected by a single network switch can be further included in a switch level group that includes all of the computer systems physically connected by a single network switch. In that case, the switch level group can be further included in a data center level group that includes all of the computer systems in a given data center. In that case, the data center level group can be further included in a universal group that includes all of the computer systems in a plurality of data centers. In such an organization, groups at each level generally have progressively slower access to copies of programs located on other computer systems within the group, with a single machine group being the fastest. Provide access, and Universal Group provides the slowest access. Such an organization contains the smallest group in which the program execution service stores a copy of a particular program to be executed and has the resource availability needed to execute that program. Since it can be searched, it can enable efficient implementation of various policy applications that guide the optimal placement of the program to be executed. Alternatively, other embodiments may not model the computer system in the program execution service by the group at all. Such an embodiment is, for example, a network switch.
As mentioned above, various embodiments can implement various policies regarding the selection of computer systems and / or groups as candidates for running the program and / or receiving a copy of the program. In many cases, various program placement policies can inevitably involve trade-offs between factors such as reliability and efficiency (eg, startup latency, network latency, or throughput). The deployment policy is the preference of the user requesting the execution of one or more programs, the number, ID, and location of the programs currently running, the number and ID of the programs currently required to be executed, and the future. Factors such as the number and identity of programs scheduled to run, the location of previously stored copies of the program, network architecture, and geographic location can be taken into account. In addition, the default enforcement of policies can be overridden or modified in some embodiments based on user requirements or other factors. For example, certain embodiments may provide a set of default policies that can be overridden by a user preference expressed in a user's request to execute one or more programs. It is possible.
In an embodiment in which a computer system managed by a program execution service spans multiple data centers, the program execution service executes multiple instances of a single program within the same data center and / or the same data. It is possible to prefer to run multiple separate program instances for the same user within the center. Such policies tend to allow such programs to utilize relatively high bandwidth intra-data center data exchange for communication between program instances. On the other hand, some embodiments disperse such program instances across multiple data centers in the case of a power outage, network outage, or other major outage that can cause the entire data center to go down. In order to ensure reliability, it may be preferred for program instances that perform little or no communication with other such program instances. Such preferences for distributing or consolidating such program instances are equally applicable at various other levels of computer system organization, such as with respect to physical subnetworks, groups, and individual computer systems. In addition, some embodiments use policies that can be used to make selections between multiple otherwise indistinguishable candidate computer systems under the program execution service deployment policy. It is possible to do. For example, one embodiment may randomly select one computer system from a set of equally good candidate computer systems, while another embodiment is the computer with the lowest resource utilization. It is possible to select systems, while different embodiments allow such computer systems to be selected in round robin order.
In addition, various embodiments may implement different policies regarding the storage of a copy of the program in the local program storage repository with respect to the execution of the program. For example, in some embodiments, a local copy of a program may always be stored on the local program storage repository prior to (or during, or after) execution on the computer system that houses the local program storage repository. It is possible. Alternatively, in other embodiments, only a few programs are stored in such a local program storage repository. Moreover, various embodiments may take different approaches if the program storage repository does not have enough capacity to store a local copy of a given program. For example, some embodiments have recently been the least used copy, the oldest copy, a random copy, a differently selected copy, a program repository of one or more other groups in a common group, and so on. Was stored in the program repository to make room for new programs, such as forcing a copy of the program that is still stored in some other related program repository, etc. Choose to evict one or more copies of the program, or remove it otherwise. In other embodiments, if a given program repository is full, no evacuation will be performed (for example, instead, daily, after a restart, etc., remove all programs from the program repository on a regular basis. By removing the program, or only if the program is unregistered from the program execution service).
In some embodiments, the program can be decomposed into multiple, and optionally, fixed size data blocks. By disassembling the program in this way, the computer system that is acquiring a copy of the program delivers the request to multiple other computer systems that store the requested program in the program repository. be able to. Since some of these other plurality of computer systems respond to requests for program blocks, the acquiring computer system can request additional program blocks from those responding computer systems. Therefore, a computer system with sufficient resources available is preferred to supply program blocks over a less responsive or unresponsive computer system.
In some embodiments, optimization is performed to improve the transfer efficiency of the program by transferring only a part of the program different from other programs that may already be stored in the local program repository. Can be done. Such an approach can be advantageous given the same program, or multiple incremental versions of different programs that share a significant portion of the code or data. For example, if a program is broken down into multiple, and possibly fixed-sized blocks, a checksum will be calculated and stored for each block when the program is first registered with the program execution service. Can be done. Later, when a program should be acquired for execution, the computer system compares the program block checksum to the checksum associated with the block of program that resides in one or more program repositories. , Only program blocks that are no longer stored can be acquired. Alternatively, some embodiments may represent the program as a collection of one or more files, such as an executable file, a data file, and a library file. In such cases, the two programs can have one or more files (eg, library files) in common, and a given computer system of the programs to be acquired for execution. You can choose to acquire only files that are different from those already stored in your computer's program repository.
Some embodiments all provide fixed size programs, while other embodiments allow programs of various sizes. Fixed-size programs can simplify the handling of programs in the context of calculating program utilization of system resources such as memory or program repositories. In embodiments that provide programs of various sizes, the utilization of fixed-size resources (such as memory or disk space) is optimized to localize the program, including various binpacking algorithms such as best fit, first fit, etc. Various algorithms are applicable that limit fragmentation when storing copies and / or when running program instances.
In addition, some embodiments provide the ability to put a copy of the program into or otherwise distribute a copy of the program into various systems of the managed computer system prior to the request to execute the program. Can be done. Some embodiments provide at least one universal program repository for storing the program when the program is first registered, but these embodiments provide the program when the program is first executed. , It can suffer long wait times because it cannot be found in any program repository that is relatively local to the computer system on which the program should be executed. If such an embodiment is configured to store a local copy of the program being executed in the local program repository, subsequent executions will incur a relatively short startup latency compared to the first execution. .. The problem of relatively long startup latency with respect to the first execution of a program can be addressed by including a copy of the program or otherwise delivering it prior to the request to execute the program. Such an embodiment can deliver one or more copies of a program to a program repository that is local to one or more data centers that provide program execution services. In such a way, when a program is requested to be run first, the program is generally relatively local to the computer system or multiple computer systems chosen to run the program ( Found in the program repository (at least in the same data center).
Further, in some embodiments, optimization can be performed in the case of simultaneous execution of multiple instances of a single program, i.e., overlapping initiations. In such situations, a copy of the program to be executed may need to be obtained by multiple separate computer systems at about the same time. If each computer system obtains a copy of the program independently from the remote program repository, then each computer system initiates the transfer of the same data over the network at the same time, resulting in overuse of the network and other resources. May be brought. Synchronize the acquisition of one or more copies of a program so that multiple computer systems make better use of system resources (eg, by minimizing unnecessary network usage) in some situations. It may be beneficial to do or otherwise order. For example, if multiple computer systems selected to run a program are part of the same group and the program copy should be obtained from one or more computer systems outside that group, the multiple computer systems. It may be beneficial for the first computer system of a computer system to first obtain a copy of the program from a computer system outside the group (and store it in a local program repository). After the first computer system obtains a copy of the program, the remaining systems of the plurality of computer systems can obtain a copy from the first computer system via a common data exchange medium for that group. ..
In addition, a variety of additional techniques are available that make efficient use of network and / or other computer resources when multiple computer systems should each obtain a copy of the program. For example, the first system of a plurality of computer systems can be selected to manage the delivery of the program to the other systems of the plurality of computer systems. If none of the computer systems have a stored copy of the program in the local program repository, the selected computer system initiates the transfer of at least a portion (eg, block) of the program from a remote location. be able to. As a portion of the program is received by the selected computer system, the selected computer system can multicast the received portion to the other systems of its plurality of computer systems. Such multicast sends fewer redundant data packets to the network connecting the multiple computer systems, so other network communication mechanisms (eg, TCP-based forwarding by each of the multiple computer systems). Compared with, it can bring about improved network usage. Alternatively, if one or more of the multiple computer systems have a stored copy of the program in the local program repository, the selected computer system will have at least the program on the other system of that multiple computer system. One with a stored copy of the program to multicast parts (eg blocks) to distribute the load of forwarding blocks and minimize the impact on other computer systems and / or parts of the network. Or it can be directed to at least some of the multiple computer systems. After such multicast-based delivery of the program, one or more of the multiple computer systems will receive An alternative communication mechanism (eg TCP) can be utilized to capture parts of the program that were not trusted (eg, for dropped network packets). An alternative delivery mechanism is a delivery request that seeks parts in a round robin manner or otherwise that distributes the load on other systems of its multiple computer systems and / or on some parts of the network. Can be included.
In some embodiments, additional techniques are available. For example, if a multicast-based delivery mechanism is used to deliver parts of a program to one group of computer systems from another computer system within that group, then groups by multicast using various techniques. External network traffic can be prevented or restricted. For example, a short lifetime can be specified for using multicast packet and / or packet addressing techniques so that the switch does not transmit multicast packets to computer systems that are not connected to it. In addition, some embodiments further / or network resources to minimize network resource usage, to minimize the load on computer systems that are not involved in the transfer or execution of a copy of the program for execution. Various policies can be implemented to provide predictable performance of and / or computer resources. For example, some embodiments limit the rate at which a computer system can transfer a copy of a program to another computer system, whether for multicast transmission or / or point-to-point transmission. It is possible. Further, in some embodiments, an intermediate network device (eg, a switch, router, etc.) may be utilized by the intermediate network device when transferring a data packet that carries a portion of a copy of the program between subnetworks. It is possible to limit the possible transfer rates and / or the percentage of network bandwidth. Such data packets are identified by intermediate network devices, for example, based on being of a particular type and / or being a particular address (eg, a multicast IP address within a certain range). It is possible. In some embodiments, a plurality of machines, such as the mechanisms described above.
In some embodiments, various techniques for transferring one or more program instances running from one or more computer systems to another one or more computer systems may also be used. .. In one aspect, this migration can reflect problems associated with the first computer system in which the program instance is running (eg, a computer system and / or a failure of network access to the computer system). In another aspect, migration is the maintenance, energy of the original computer system that executes the program instance for high priority program execution or by integrating the execution of the program instance on a limited number of computer systems. You can deal with other program instances that should be running on the first computer system, such as allowing them to be shut down for reasons such as savings. As one particular example, if one or more program instances running on a computer system require more resources available from that computer system, then one or more of these program instances. May need to be transferred to one or more other computer systems with additional resources. One or more computer systems have less resources than expected, one or more computer systems use more resources than expected (or allowed), and so on. , Or an implementation in which the available resources of one or more computer systems are deliberately overcommitted compared to the possible resource needs of one or more reserved or running program instances. In the form, excessive use of available resources can occur. For example, if the expected resource needs of a program instance are within the available resources, then the biggest resource needs are: May exceed available resources. Also, if the actual resources required to execute a program instance exceed the available resources, excessive use of the available resources may occur. Program migration is like transferring a copy of a program stored locally on the first computer system to the target destination computer system, and / or on the target destination computer system of the program running on the first computer system. It can be executed in various ways, such as starting to execute a new instance. The migration allows the current execution state information to be transferred to the new running program instance, and / or allows other coordination between the first program instance and the new program instance, and so on. If possible, it can be done before the first executed program instance is terminated.
In some embodiments, the program execution service can be provided to a plurality of customers in exchange for a fee. In such a situation, the customer may register the program with the program execution service or otherwise provide it and request the execution of such a program in exchange for a fee. One or so that the customer purchases access to various configurations of program execution service resources (eg, network bandwidth, memory, storage, processor) on an hourly basis (eg, minutes, hours, days, etc.) Premium services (eg, start running a premium customer's program prior to a non-premium customer's program, to purchase access to multiple given virtual or physical hardware configurations, premium Prior to the customer's program, it provides preferred execution, such as eviction of programs belonging to non-premium customers, resulting in a preferred program repository placement, etc.) for a specified period of time per instance execution. Various billing models are available, such as purchasing the ability to run a program instance.
As mentioned above, some embodiments are capable of using a virtual computer system, in which case the program to be executed by the program execution service can include the entire virtual computer image. .. In such embodiments, the program to be executed can include the entire operating system, file system, and / or other data, and optionally one or more user-level processes. In other embodiments, the program to be executed may include one or more other types of executable files that interact to provide some functionality. In yet other embodiments, the program to be executed is native on the computer system provided or indirectly using a hardware abstraction implemented by a virtual computer system, interpreter, or other software. It can include a physical or logical set of executable instructions and data. More generally, in some embodiments, the program to be executed includes one or more application programs, application frameworks, libraries, archives, class files, scripts, configuration files, data files, and the like. You may.
An embodiment that utilizes a combination of a mutual communication system manager module and a machine manager module to manage the execution of a program within a program execution service has been described, but other responsibilities between the various program execution service modules. Implementation and allocation are also planned. For example, in some embodiments, a single module, or single component, executes a program on some, or all, or machines of a managed physical computer system or virtual machine. It is possible to take charge of managing. For example, the program can be executed directly on the target computer system by various remote execution techniques (eg rexec, rsh, etc.).
It will also be appreciated by those skilled in the art that although the exemplary embodiments described above have been used in the context of the data center used to provide program execution services, other implementation scenarios are possible. .. For example, the facilities described are used in the context of an organization-wide intranet operated by a business or other institution (eg, a university) for the benefit of employees and / or other members of that business or institution. May be done. Alternatively, the techniques described are distributed, including nodes that are individually managed and operated by various third parties for the purpose of performing large-scale (eg, scientific) computing tasks in a distributed manner. It may be used by a computer system.
From the above, although specific embodiments have been described herein for illustrative purposes, it will be appreciated that various modifications can be made without departing from the spirit and scope of the present invention. Therefore, the present invention is not limited except by the appended claims and the elements described in the claims. Further, although some aspects of the invention are presented in some claims, the inventors contemplate various aspects of the invention in any claim. For example, only some aspects of the invention may be described as being realized on a computer-readable medium, but other aspects as well are possible.
8 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8
Every citation, both ways
| Reference | Relation | Cited during |
|---|---|---|
| JPN6011033846; Eun-Kyu Byun;Jae-Wan Jang;Wook Jung;Jin-Soo Kim: 'A dynamic grid services deployment mechanism for on-demand resource provisioning' Proc. of 5th Int. Symp on Cluster Computing and the Grid Vol.2, 20050509, P.863-870 | Non-patent | Search report |
| JPN6012033211; 前園ほか: 'データグリッドにおけるデータ転送を考慮したローカルスケジューリング手法' 電子情報通信学会技術研究報告 第104巻,第692号, 200502, P.435-440, 社団法人電子情報通信学会 | Non-patent | Search report |
| JPN6012033213; 滝澤真一朗ほか: 'グリッド上のスケーラブルな並列レプリケーションフレームワーク' 情報処理学会研究報告 第2004巻,第81号, 20040801, P.247-252, 社団法人情報処理学会 Information Processing Socie | Non-patent | Search report |
| JPN6012033215; 渡邊ほか: 'GridRPCシステムにおけるリモートプログラムシッピング機構' 情報処理学会研究報告 第2003巻,第102号, 20031016, P.73-78, 社団法人情報処理学会 | Non-patent | Search report |
| JPN6012033216; 山形ほか: 'グリッド上における仮想計算機を用いたジョブ実行環境構築システムの高速化' 情報処理学会研究報告 第2006巻,第20号, 20060228, P.127-132, 社団法人情報処理学会 | Non-patent | Search report |
| CSNG200600628084; 前園ほか: 'データグリッドにおけるデータ転送を考慮したローカルスケジューリング手法' 電子情報通信学会技術研究報告 第104巻,第692号, 200502, P.435-440, 社団法人電子情報通信学会 | Non-patent | Examiner |
| CSNG200500584040; 滝澤真一朗ほか: 'グリッド上のスケーラブルな並列レプリケーションフレームワーク' 情報処理学会研究報告 第2004巻,第81号, 20040801, P.247-252, 社団法人情報処理学会 Information Processing Socie | Non-patent | Examiner |
| CSNG200401901013; 渡邊ほか: 'GridRPCシステムにおけるリモートプログラムシッピング機構' 情報処理学会研究報告 第2003巻,第102号, 20031016, P.73-78, 社団法人情報処理学会 | Non-patent | Examiner |
| CSNG200600480016; 山形ほか: 'グリッド上における仮想計算機を用いたジョブ実行環境構築システムの高速化' 情報処理学会研究報告 第2006巻,第20号, 20060228, P.127-132, 社団法人情報処理学会 | Non-patent | Examiner |
| JPN6011033846; Eun-Kyu Byun;Jae-Wan Jang;Wook Jung;Jin-Soo Kim: 'A dynamic grid services deployment mechanism for on-demand resource provisioning' Proc. of 5th Int. Symp on Cluster Computing and the Grid Vol.2, 20050509, P.863-870 | Non-patent | Examiner |
| JPN6012033211; 前園ほか: 'データグリッドにおけるデータ転送を考慮したローカルスケジューリング手法' 電子情報通信学会技術研究報告 第104巻,第692号, 200502, P.435-440, 社団法人電子情報通信学会 | Non-patent | Examiner |
| JPN6012033213; 滝澤真一朗ほか: 'グリッド上のスケーラブルな並列レプリケーションフレームワーク' 情報処理学会研究報告 第2004巻,第81号, 20040801, P.247-252, 社団法人情報処理学会 Information Processing Socie | Non-patent | Examiner |
| JPN6012033215; 渡邊ほか: 'GridRPCシステムにおけるリモートプログラムシッピング機構' 情報処理学会研究報告 第2003巻,第102号, 20031016, P.73-78, 社団法人情報処理学会 | Non-patent | Examiner |
| JPN6012033216; 山形ほか: 'グリッド上における仮想計算機を用いたジョブ実行環境構築システムの高速化' 情報処理学会研究報告 第2006巻,第20号, 20060228, P.127-132, 社団法人情報処理学会 | Non-patent | Examiner |
80 members in 7 offices
Priority claims2
| Document | Office | Kind | Date |
|---|---|---|---|
| 11395463 | United States of America | – | |
| 39546306 | United States of America | A |
Members80
| Document | Office | Kind | |
|---|---|---|---|
| US4819312A | United States of America | A | |
| CA1314156C | Canada | C | |
| US2007239987A1 | United States of America | A1 | |
| US2007240160A1 | United States of America | A1 | |
| CA2646135A1 | Canada | A1 | |
| CA2646157A1 | Canada | A1 | |
| CA2875381A1 | Canada | A1 | |
| CA2920179A1 | Canada | A1 | |
| WO2007126835A2 | World Intellectual Property Organization (WIPO) | A2 | |
| WO2007126837A2 | World Intellectual Property Organization (WIPO) | A2 | |
| US2008059557A1 | United States of America | A1 | |
| WO2007126835A3 | World Intellectual Property Organization (WIPO) | A3 | |
| WO2007126837A3 | World Intellectual Property Organization (WIPO) | A3 | |
| EP2008167A2 | European Patent Office (EPO) | A2 | |
| EP2008407A2 | European Patent Office (EPO) | A2 | |
| CA2697540A1 | Canada | A1 | |
| WO2009032471A1 | World Intellectual Property Organization (WIPO) | A1 | |
| CN101460907A | China | A | |
| CN101461190A | China | A | |
| JP2009532771A | Japan | A | |
| JP2009532944A | Japan | A | |
| EP2186012A1 | European Patent Office (EPO) | A1 | |
| US7792944B2 | United States of America | B2 | |
| US7801128B2 | United States of America | B2 | |
| US2010312871A1 | United States of America | A1 | |
| US2010318645A1 | United States of America | A1 | |
| EP2008167A4 | European Patent Office (EPO) | A4 | |
| US8010651B2 | United States of America | B2 | |
| US8190682B2 | United States of America | B2 | |
| JP2012142012AThis record | Japan | A | |
| JP5006925B2 | Japan | B2 | |
| EP2186012A4 | European Patent Office (EPO) | A4 | |
| CA2646135C | Canada | C | |
| JP5174006B2 | Japan | B2 | |
| CN101461190B | China | B | |
| US8509231B2 | United States of America | B2 | |
| US2013283176A1 | United States of America | A1 | |
| JP5335948B2 | Japan | B2 | |
| JP2013229040A | Japan | A | |
| US2013298191A1 | United States of America | A1 | |
| CN101460907B | China | B | |
| CN103713951A | China | A | |
| EP2008407A4 | European Patent Office (EPO) | A4 | |
| IN2671KON2014A | India | A | |
| JP2015144020A | Japan | A | |
| JP5789640B2 | Japan | B2 | |
| US9253211B2 | United States of America | B2 | |
| US2016142254A1 | United States of America | A1 | |
| US9426181B2 | United States of America | B2 | |
| US2016359919A1 | United States of America | A1 | |
| US9621593B2 | United States of America | B2 | |
| JP6126157B2 | Japan | B2 | |
| CA2875381C | Canada | C | |
| CA2646157C | Canada | C | |
| US2017212780A1 | United States of America | A1 | |
| CA2697540C | Canada | C | |
| US9794294B2 | United States of America | B2 | |
| US2018048675A1 | United States of America | A1 | |
| EP2186012B1 | European Patent Office (EPO) | B1 | |
| CN103713951B | China | B | |
| CA2920179C | Canada | C | |
| US10348770B2 | United States of America | B2 | |
| US10367850B2 | United States of America | B2 | |
| US2020007587A1 | United States of America | A1 | |
| US2020021617A1 | United States of America | A1 | |
| EP2008167B1 | European Patent Office (EPO) | B1 | |
| EP2008407B1 | European Patent Office (EPO) | B1 | |
| US10764331B2 | United States of America | B2 | |
| EP3702919A1 | European Patent Office (EPO) | A1 | |
| US10791149B2 | United States of America | B2 | |
| EP3731465A1 | European Patent Office (EPO) | A1 | |
| US2021014279A1 | United States of America | A1 | |
| US2021051181A1 | United States of America | A1 | |
| US11451589B2 | United States of America | B2 | |
| US11539753B2 | United States of America | B2 | |
| US2023084547A1 | United States of America | A1 | |
| US2023208884A1 | United States of America | A1 | |
| EP3731465B1 | European Patent Office (EPO) | B1 | |
| US11997143B2 | United States of America | B2 | |
| US12003548B2 | United States of America | B2 |
26 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Request for written amendment filedJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Written permission of extension of timeJAPANESE INTERMEDIATE CODE: A602A602 | A602 | |
| Written request for extension of timeJAPANESE INTERMEDIATE CODE: A601A601 | A601 | |
| Written permission of extension of timeJAPANESE INTERMEDIATE CODE: A602A602 | A602 | |
| Written request for extension of timeJAPANESE INTERMEDIATE CODE: A601A601 | A601 | |
| Written permission of extension of timeJAPANESE INTERMEDIATE CODE: A602A602 | A602 | |
| Written request for extension of timeJAPANESE INTERMEDIATE CODE: A601A601 | A601 | |
| Notification of reasons for refusalJAPANESE INTERMEDIATE CODE: A131A131 | A131 | |
| Report on accelerated examinationJAPANESE INTERMEDIATE CODE: A971005A975 | A975 | |
| Request for written amendment filedJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Explanation of circumstances concerning accelerated examinationJAPANESE INTERMEDIATE CODE: A871A871 | A871 |
Numbers
- Publication
- 2012142012
- Application
- 52264
Titles2
- Japanese
- 複数のコンピュータシステムによるプログラムの実行を管理するシステム
- English
- A system that manages the execution of programs by multiple computer systems
Classification
- CPC, 2
- G06F8/60
- G06F9/5055
- IPC, 1
- G06F9 50