An infrastructure for parallel programming of clusters of machines
Abstract
INFRASTRUCTURE FOR PARALLEL PROGRAMMING OF GROUPS OF MACHINES. GridBatch provides an infrastructure support structure that hides the complexities and burdens of programming development and application logic that implements parallel computations of programmer details. A programmer can use GridBatch to implement parallelized computational operations that minimize network bandwidth requirements, and efficiently divide and coordinate processing with putational in a multi-processor configuration. GridBatch provides an efficient and lightweight approach to quickly build parallelized applications using economically viable multiprocessor configurations that achieve the highest performance results.
Term
Projected expiry 30 September 2028.
- Priority
- Filed
- Granted
- Today
- Projected expiry
21 claims: 1 independent, 20 dependent
- 1REIVINDICAÇÕES 1. Produto compreendendo:um meio legível por máquina;primeira lógica de operador armazenada no meio e operável para: implementar uma primeira operação de processamento de dados em paralelo sobre múltiplos nós de processamento, a primeira operação de processamento customizada com uma primeira função definida pelo usuário executada nos múltiplos nós de processamento;e segunda lógica de operador armazenada no meio e operável para: implementar uma segunda operação de processamento de dados em paralelo por todos os múltiplos nós de processamento, a segunda operação de processamento de dados customizada com uma segunda função definida pelo usuário executada nos múltiplos nós de processamento.
- 2Produto de acordo com a reivindicação 1, compreendendo adicionalmente:lógica do administrador do sistema de arquivo armazenada no / meio 7 , e operável para: designar blocos de vetor de um primeiro vetor entre os múltiplos nós de processamento de acordo com uma função unidirecional definida pe-lo-usuári©:-------- - ------------ - —
- 3Produto de acordo com a reivindicação 2, onde a lógica do administrador do sistema de arquivo é adicionalmente operável para:fornecer informação de localização de nó de bloco do vetor para os blocos de vetor para um programador de trabalho.
- 4Produto de acordo com a reivindicação 2, onde a lógica do administrador do sistema de arquivo é adicionalmente operável para:rearranjar os blocos do vetor.
- 5Produto de acordo com a reivindicação 2, onde a lógica do administrador do sistema de arquivo é adicionalmente operável para:manter um mapeamento de IDs de bloco para os múltiplos nós de processamento que identifica cada designação de nó de dados da ID do bloco;e rearranjar os blocos do vetor quando o mapeamento muda.
- 6Produto de acordo com a reivindicação 1, onde:5 a primeira ou segunda lógica do operador compreende lógica do operador de ligação;a primeira ou segunda função definida pelo usuário compreende uma função de ligação definida pelo usuário;e onde a lógica do operador de ligação é operável para invocar a 10 função de ligação definida pelo usuário em um primeiro registro combinado em um primeiro vetor e um segundo registro combinado em um segundo vetor distribuída entre os múltiplos nós de processamento quando o campo de índice de ligação presente no primeiro vetor e no segundo vetor combina para o primeiro registro combinado e o segundo registro combinado, para 15 obter um resultado de ligação.
- 7Produto de acordo com a reivindicação 6, compreendendo adicionalmente:lógica de nó máster armazenada no meio e operável para: receber uma chamada de função de ligação;e 20 iniciar geração de tarefas de ligação localmente entre os múltiplos nós de processamento, cada tarefa de ligação operável para seletivamente iniciar a execução da função de ligação definida pelo usuário.
- 8Produto de acordo com a reivindicação 1, onde:a primeira ou segunda lógica do operador compreende lógica de 25 operador recursivo;a primeira ou segunda função definida pelo usuário compreende uma função recursiva definida pelo usuário;e onde a lógica do operador recursivo é operável para invocar a função recursiva definida pelo usuário iniciando sobre os blocos de vetor 30 localmente nos múltiplos nós de processamento para produzir resultados intermediários, comunicar um subconjunto dos resultados intermediários a um subconjunto dos múltiplos nós de processamento, e iterar: invocar a função recursiva definida pelo usuário nos resultados intermediários para produzir progressivamente menos resultados intermediários;e comunicar a um subconjunto dos progressivamente menos resultados intermediários para um subconjunto dos progressivamente menor dod múltiplos nós de processamento;até um resultado recursivo final ser obtido sobre o primeiro vetor em um nó final no primeiro conjunto de nós.
- 9Produto de acordo com a reivindicação 8, compreendendo adicionalmente:lógica de nó máster armazenada no meio e operável para: receber uma chamada de função recursiva;e iniciar geração de tarefas de operação recursiva localmente entre os múltiplos nós de processamento, cada tarefa de operação recursiva operável para seletivamente iniciar a execução da função recursiva definida pelo usuário para os blocos de vetor.
- 10Produto de acordo com a reivindicação 1, onde:a primeira ou segunda lógica do operador compreende lógica do operador de convolução;a primeira ou segunda função definida pelo usuário compreende uma função de convolução definida pelo usuário;e ----- onde a lógica do operador de convolução é operável para invocar a função de convolução definida pelo usuário para cada registro em um primeiro vetor em cada registro em um segundo vetor, para obter um resultado de função de convolução.
- 11Produto de acordo com a reivindicação 10, compreendendo adicionalmente:lógica de nó máster armazenada no meio e operável para: receber uma chamada de função de convolução;e iniciar geração de tarefas de operação de convolução localmente entre os múltiplos nós de processamento, cada tarefa de operação de convolução operável para seletivamente iniciar execução da função de convolução definida pelo usuário.
- 12Produto de acordo com a reivindicação 1, onde:a primeira ou segunda lógica do operador compreende distribuir lógica do operador;5 a primeira ou segunda função definida pelo usuário compreende uma função de divisão definida pelo usuário;e onde a lógica do operador de distribuição é operável para redistribuir, de acordo com a função de divisão definida pelo usuário, um primeiro vetor previamente distribuído como primeiros blocos do vetor entre os múlti10 pios nós de processamento, para obter blocos de vetor redistribuídos do primeiro vetor redistribuído entre os múltiplos nós de processamento.
- 13Produto de acordo com a reivindicação 1, onde:a primeira ou segunda lógica do operador compreende lógica do operador;
- 1415 a primeira ou segunda função definida pelo usuário compreende uma função de mapa definida pelo usuário; e onde a lógica do operador de mapa é operável para aplicar a função de mapa definida pelo usuário para registros de um vetor distribuído entre os múltiplos nós de processamento. 20 14. Método para processamento de dados em paralelo compreendendo:------ iniciar-a execução-de uma-primeira-operação de processamento de dados em paralelo sobre múltiplos nós de processamento, a primeira operação de processamento de dados customizada com uma primeira função 25 definida pelo usuário executada nos múltiplos nós de processamento;e iniciar a execução de uma segunda operação de processamento de dados em paralelo sobre os múltiplos nós de processamento, a segunda operação de processamento de dados customizada com uma função definida pelo usuário executada nos múltiplos nós de processamento. 30 15. Método de acordo com a reivindicação 14, compreendendo adicionalmente: designar blocos de vetor de um primeiro vetor entre os múltiplos nós de processamento de acordo com uma função unidirecional definida pelo usuário.
- 1516. Método de acordo com a reivindicação 15, compreendendo adicionalmente:fornecer informação da localização de nó do bloco do vetor para os blocos do vetor a um programador de trabalho.
- 1617. Método de acordo com a reivindicação 15, compreendendo adicionalmente:rearranjar os blocos de vetor.
- 1718. Método de acordo com a reivindicação 14, onde:a primeira ou segunda lógica do operador compreende lógica do operador de ligação;a primeira ou segunda função definida pelo usuário compreende uma função de ligação definida pelo usuário;e onde a lógica do operador de ligação invoca a função de ligação definida pelo usuário em um primeiro registro combinado em um primeiro vetor e um segundo registro combinado em um segundo vetor distribuído entre os múltiplos nós de processamento quando um campo de índice de ligação presente no primeiro e no segundo vetor combina com o primeiro registro de combinação e o segundo registro de combinação, para obter um resultado de ligação.
- 1819. Método de-acordo-com a reivindicação-18, compreendendo adicionalmente:receber uma chamada de função de ligação;e iniciar rearranjo de tarefas de ligação localmente entre os múltiplos nós de processamento, cada tarefa de ligação operável para seletivamente iniciar a execução da função de ligação definida pelo usuário.
- 1920. Método de acordo com a reivindicação 14, onde:a primeira ou segunda lógica do operador compreende lógica do operador recursivo;a primeira ou segunda função definida pelo usuário compreende uma função recursiva definida pelo usuário;e onde a lógica do operador recursivo invoca a função recursiva definida pelo usuário iniciando sobre os blocos do vetor localmente nos múltiplos nós de processamento para produzir resultados intermediários, comunica um subconjunto dos resultados intermediários para um subconjunto dos 5 múltiplos nós de processamento, e itera: invocar a função do recurso definido pelo usuário nos resultados intermediários para produzir progressivamente menos resultados intermediários;e comunicar um subconjunto dos progressivamente menos resul10 tados intermediários para um subconjunto progressivamente menor dos múltiplos nós de processamento;até um resultado recursivo final ser obtido sobre o primeiro vetor em um nó final no primeiro conjunto de nós.
- 2021. Método de acordo com a reivindicação 20, compreendendo 15 adicionalmente:receber uma chamada de função recursiva;e iniciar a geração de tarefas de operação recursiva localmente entre os múltiplos nós de processamento, cada tarefa de operação recursiva operável para seletivamente iniciar a execução da função recursiva definida 20 pelo usuário para os blocos de vetor.
- 2122. Método de acordo com a reivindicação 14, onde:-a-primeira-ou-a-segunda-lógiea-do operador compreende lógica do operação de convolução;a primeira ou a segunda função definida pelo usuário compreen25 de uma função de convolução definida pelo usuário;e onde a lógica do operador de convolução invoca a função de convolução definida pelo usuário para cada registro em um primeiro vetor _em todo registro em um segundo vetor, para obter um resultado da função de convolução. 30 23. Método de acordo com a reivindicação 22, compreendendo adicionalmente: receber uma chamada da função de convolução;e iniciar geração de tarefas de operação de convolução localmente entre os múltiplos nós de processamento, cada tarefa de operação de convolução operável para seletivamente iniciar a execução da função de convolução definida pelo usuário. 5 24. Método de acordo com a reivindicação 14, onde: a primeira ou a segunda lógica do operador compreende distribuir lógica do operador;a primeira ou a segunda função definida pelo usuário compreende uma função de divisão definida pelo usuário;e 10 onde a lógica do operador de distribuição redistribui, de acordo com a função de divisão definida pelo usuário, um primeiro vetor previamente distribuído como blocos do primeiro vetor entre os múltiplos nós de processamento, para obter blocos de vetor redistribuídos do primeiro vetor redistribuído entre os múltiplos nós de processamento. 15 25. Método de acordo com a reivindicação 14, onde: a primeira ou a segunda lógica do operador compreende lógica do operador de mapa;a primeira ou a segunda função definida pelo usuário compreende uma função definida pelo usuário;e 20 onde a lógica do operador do mapa aplica a função do mapa definida pelo usuário para registros de um vetor distribuído entre os múltiplos nós de processamento.........- - —------ _ -------1/12 Nó escravo 120 100 Interface de comunicações 113 112 Processador Armazenagem Ύ“ Cluster de GridBatch Nó máster 2/12 116 Interface de comunicações 211 Processador 210 Memória 215 Lógica de administração 222 do sistema do arquivo Lógica de programador 230 de trabalho Identificador do 272 primeiro vetor Identificador do 274 segundo vetor Função definida 276 pelo usuário Identificador do vetor 280 de resultados índice do vetor Lógica do nó máster 260 Biblioteca de software 262 de GridBatch Solicitação de tarefa 244 Tarefa
Independent claims21
132 paragraphs in 1 section, as filed
(54) Title: INFRASTRUCTURE FOR PARALLEL PROGRAMMING OF GROUPS OF MACHINES (30) Unionist Priority: 01/10/2007 us 11 / 906,293 (73) Owner (s): Accenture Global Services Gmbh (72) Inventor (s): Huan Liu (57) Summary: infrastructure for PARALLEL programming of GROUPS OF MACHINES. GridBatch provides an infrastructure support structure that hides the complexities and burdens of programming development and application logic that implements parallel computations of programmer details. A programmer can use GridBatch to implement parallelized computational operations that minimize network bandwidth requirements, and efficiently divide and coordinate processing with putational in a multi-processor configuration. GridBatch provides an efficient and lightweight approach to quickly build parallelized applications using economically viable multiprocessor configurations that achieve the highest performance results.
<td colspan="2">Hollow knot i2 £</td>
<td>r |</td><td> “1</td>
<td>t</td><td>Ui I</td>
<td colspan="2">Jlâ Memory</td>
<td colspan="2"> 1“*- “1</td>
<td colspan="2">- | Slave chores 1 £ 5 |</td>
<td>I No logic</td><td>16Q mo |</td>
<td colspan="2"></td>
<td>L (Store</td><td>ΞΖϋ</td>
<td colspan="2"></td>
<img file="BRPI0805054A2_D0001.tif" />
ΡΙ0805054-6
Descriptive Report of the Invention Patent for INFRASTRUCTURE FOR PARALLEL PROGRAMMING OF GROUPS OF MACHINES. <sup>Λ</sup>
Background of the Invention
1. Technical Field: This description refers to a system and method for parallelizing applications through the use of a software library of operators designed to implement paralyzed computing plans in detail. In particular, this description refers to an efficient and cost-effective way to implement parallelized applications.
2. Background of the Invention
There is currently a wide disparity between the number of data organizations needed to process at any given time and the computing capacity available to the organization using single CPU (uniprocessor) systems. Today, organizations use applications that process terabytes and even petabytes of data in order to derive valid information and business insight. Unfortunately, many applications typically operate uniprocessor machines sequentially, and require hours and even days of computing time to produce usable results. The gap between the amount of data that organizations can process and the computational performance of uniprocessors available to organizations continues to widen. The amount of data collected and processed by organizations continues to grow exponentially. Organizations must address enterprise database growth rates of around 125% year on year or the equivalent to double in size every 10 months. The volume of data for other data-rich industries also continues to grow exponentially. For example, Astronomy has a data doubling rate every 12 months, every 9 months for Biosequences, and every 6 months for Functional Genomics.
Although storage capacity continues to grow at an exponential rate, the speed of uniprocessors no longer grows exponentially. In this way, even though organizations may have the ability to continue to increase data storage capacity, the computational performance of uniprocessor configurations may no longer keep pace. Organizations must identify a technical solution to address diverging trends in storage capacity and uniprocessor performance.
In order to process large amounts of data, applications require large amounts of computing power and high I / O performance. Programmers face the technical challenges of yet more efficient identification to divide computational processing and coordinate multiple CPUs through computing to address the widening gap between demand and supply of computing capacity. Given the limited reality of network bandwidth availability, programmers also face the technical challenge of addressing the high bandwidth requirements needed to distribute vast amounts of data to multiple CPUs running parallel processing computations. Simply introducing an additional machine to a processing (configuration) pool does not increase the total bandwidth on the configuration network. Although the I / O bandwidth of the local disk may increase as a result. A network topology can be represented as a tree that has many branches that represent the network segments and leaves that represent the processors. In this way, a single bottleneck across any one network segment can determine the total network capacity and bandwidth of a configuration. In order to scale the bandwidth, the efficient use of increases in the l / O bandwidth of the local disk should be leveraged.
The extraordinary technical challenges associated with computational parallelization operations include parallel programming complexity, proper development and testing tools, scaling limits for network bandwidth, diverging storage capacity trends and uniprocessor performances, and division efficient computational processing and coordination in multiprocessor configurations.
There has been a need for a long time for a system and method that economically and efficiently implements parallel computing solutions and effectively frees the burden of developing complex parallel programs by programmers.
summary
The GridBatch system provides an infrastructure context that programmers can use to easily convert a high-level project into a parallelized computational implementation. The programmer analyzes the potential for parallelization of computations in an application, decomposes the computations into discrete components and considers a data division plan to obtain the highest performance. GridBatch implements the detailed parallelized computational plan, developed by the programmer without requiring the programmer to create low-level logic to perform the computations. GridBatch provides a library of operators (a primitive for data set management) as building blocks for implementing parallelization. GridBatch hides all the complexity associated with parallel programming in the GridBatch library so that the programmer only needs to understand how to apply operators to correctly implement parallelization.
Although GridBatch can support many types of applications, GridBatch -provides a particular benefit for programmers focused on deploying analytical applications, because of the unique characteristics of analytical applications and computational operators used by analytical applications. Programmers often write analytical applications to collect statistics from a large set of data, just as a particular event often occurs. The computational requirements of analytical applications often involve correlation data from two or more different sets of data (for example, the computational demands imposed by a table link expressed in an SLQ statement).
GridBatch leverages data localization techniques to efficiently manage disk I / O and effectively scale system bandwidth requirements. In other words, GridBatch divides computational processing and coordinates computing across multiple processors so that processors perform computations on local data. GridBatch minimizes the amount of data transmitted to multiple processors to perform parallel processing computations.
GridBatch solves technical problems with computational parallelization operations by hiding complexities of parallel programming, leveraging localized data to minimize network bandwidth requirements, and managing the division of computational processing and coordination between multiprocessor configurations.
Other systems, methods, and features of the invention will be, or will become apparent to, those skilled in the art upon examination of the figures and the detailed description below. Such additional systems, methods, characteristics and advantages are intended to be included in this description, are within the scope of the invention, and are protected by the appended claims.
Brief Description of Drawings
The description can be better understood with reference to the following drawings and description. The components in the figures are not necessarily-to-scale-in-emphasis rather than being placed by way of illustration of the principles of the invention. In addition, in the figures, similar reference numbers designate corresponding parts or elements across all different views.
Figure 1 illustrates the configuration of the GridBatch system.
Figure 2 shows an example of a Master Node.
Figure 3 illustrates the configuration of the GridBatch system during the processing of a call from the dispatch function.
Figure 4 shows the configuration of the GridBatch system during the processing of a call function call.
Figure 5 shows the configuration of the GridBatch system during the processing of a convolution function call.
Figure 6 illustrates the configuration of the GridBatch system during the processing of a recursive function call.
Figure 7 illustrates what logical flow the GridBatch system configuration can achieve to perform the deploy operator.
Figure 8 shows what logic flow the GridBatch system configuration can obtain to run the link operator.
Figure 9 shows what logic flow the GridBatch system configuration can obtain to perform the convolution operator.
Figure 10 shows what logic flow the GridBatch system configuration can obtain to perform the recursive function operator.
Figure 11 illustrates the configuration of the GridBatch system during the processing of a map function call.
Figure 12 shows what logic flow GridBatch 100 can obtain to run the map operator.
Detailed Description
Previous research in parallel computing focused on automatically detecting parallelism in a sequential application. For example, engineers developed techniques in computer architecture, such as out-of-order buffers, designed to detect dependencies between instructions and parallel schema-independent instructions. Such techniques only examine code fragments, encoded in a sequential programming language and cannot exploit the level of application of parallelism. In this way, such techniques limit the amount of parallelism that can be explored.
A large class of applications, in particular data intensive batch file applications, have obvious parallelism at the data level. However, there are several technical challenges for implementing parallel applications. Programmers must address non-trivial issues regarding communications, coordination and synchronization between machines and processors when programmers design a parallelized application. In stark contrast to sequential programs, programmers must anticipate all possible interactions between all machines in the configuration of a parallel program, given the inherent asynchronous nature of parallel programs. Also, efficient debugging tools for parallelized application and configuration development do not exist. For example, stepping through some code can be difficult to perform in an environment where the configuration has many strands operating on many machines. Also, because of the complex interactions that result in parallelized applications, programmers identify many flaws seen as transient in nature and difficult to reproduce. The technical challenges faced by programmers who implement parallelized applications directly translate into greater cost development and development of longer cycles. In addition, programmers are often unable to migrate or replicate a parallelized solution to other implementations.
Programmers recognize database systems that are also suitable for analytical applications. Unfortunately, database systems do not scale for large data sets for at least two reasons. First, database systems feature a high level of SQL (Structured Query Language) in order to hide the details of the implementation. Although SQL can be relatively easy to use, the nature of such a high-level language forces -usu-users to express computations-in-a-way that results in processing that performs inefficiently from a parallelization perspective. Unlike programming in a low level language (for example, C ++) where the parallelized process only reads a set of data since, the same processing expressed in SQL can result in several readings being performed. Even though the techniques have been developed to automatically optimize query processing, the performance performed using a low-level language to implement parallel computing still exceeds the performance of the higher-level language such as SQL. Second, the l / O architecture of database systems limits the scalability of distributed parallel implementations because the databases assume that data access is via a common logical storage unit on the network, or through a distributed file system. or hardware SAN (network storage area). Databases do not leverage logic for physical data mapping and, therefore, do not take advantage of data locality or physical data location. Even though sophisticated hiding mechanisms exist, databases often access data by crossing the network unnecessarily and consuming precious network bandwidth.
Analytical applications differ from web applications in several ways. Analytical applications typically process structured data, while web applications often deal with unstructured data. Analytical applications often require cross-reference information from different sources (for example, different database tables). Analytical applications typically focus on statistics much less than applications on the web. For example, a word count application would require statistics for all words in a vocabulary, where an analytical application may only be interested in the number of products sold.
GridBatch provides fundamental operators that can be used for analytical or other applications. A detailed parallelized application implementation can be expressed as a combination of basic operators provided by GridBatch. GridBatch saves the programmer considerable time in terms of implementation and debugging because GridBatch addresses the parallel aspects of programming for the programmer. Using GridBatch, the programmer determines the desired combination of operators, the sequence of operators, and the minimum programming to develop each operator.
Although specific GridBatch components are described, methods and systems, and articles of manufacture consistent with GridBatch may include additional or different components. For example, a processor can be implemented as a microprocessor, microcontroller, application specific integrated circuit (ASIC), discrete logic, or a combination of other types of circuits or logic. Similarly, memories can be DRAM, SRAM, Flash or any other type of memory. The logic that implements the processing and programs described below can be stored (for example, as computer executable instructions) in a computer-readable medium such as an optical or magnetic disk or other memory. Alternatively or additionally, the logic can be performed on an electromagnetic or optical signal that can be transmitted between entities. Signals, data, databases, tables and other data structures can be separately stored and managed, can be incorporated into a single memory or database, can be distributed, or can be logically and physically organized in many different ways. The programs can be part of a single program, separate programs, or distributed through different memories or processors. Furthermore, the programs, or any portion of the programs, can instead be implemented in hardware.
An example is described below where a web-based retailer sells computer equipment such as PCs and printers. The retailer uses several tables requiring terabytes of storage to track volumes of data and information that can be used to derive analytical information using several tables including: transaction table, -customer table -and-distributor table.-The table - Transaction stores the records for the product id of each item sold and the customer id of the buyer. The customer table stores customer information for every customer, and the distributor table stores information about every distributor doing business with the retailer. The retailer can use GridBatch to analyze many analytics, some of the analytics include simple counting statistics (for example, how much of a particular product has been sold and identifies the top 10 customers that produce income). The retailer can use GridBatch to analyze more complicated analytics that involve multiple tables and complex computations. For example, the retailer can use GridBatch to determine the number of customers located in geographic proximity for a retailer's distribution facility in order to measure the efficiency of the distribution network.
The GridBatch infrastructure operates on a cluster of processing nodes (nodes). Two software components operate in the GridBatch cluster environment, named the file system administrator and the job scheduler. The file system administrator manages files and stores files across all computing nodes in the cluster. The file system administrator can segment a large file into smaller blocks and store each block on separate nodes. Among all nodes in the cluster, GridBatch can designate, for example, one node to serve as the name node and all other nodes serve as data nodes.
A data node maintains a block from a large file. In an implementation, depending on the number of nodes in the cluster and other configuration considerations, a data node can maintain more than one block in a large file. A data node responds to customer requirements to read from and write to blocks assigned to the data node. The name node maintains the namespace for the file system. The name node maintains the mapping of a large file to the block list, the data nodes assigned to each block, and the physical and logical location of each data node. The name node also responds to queries from customers who request one-file allocation and allocates blocks of large files to data nodes. In an implementation, GridBatch references the nodes by the IP addresses of the nodes, so that GridBatch can access the nodes directly. The master node also maintains a physical network topology that keeps track of which nodes are directly connected. The physical topology of the network can be manually populated by an administrator and / or discovered using an automated topology discovery algorithm. The network topology information can improve the performance of the recursive functions operator by indicating slave nodes in the immediate vicinity where intermediate results can be sent and / or retrieved in order to reduce the consumption of network bandwidth. A brief description of the topology and its use in facilitating the operator's execution of recursive functions will be discussed below.
The GridBatch file system distributes large files across many nodes and tells the job programmer the location of each block so that the job programmer can schedule tasks on the nodes that host the blocks to be processed. GridBatch targets large-scale data analysis problems, such as data storage, where a large amount of structured data needs to be structured. A file typically stores a large collection of data records that have an identical schema (for example, object owner, or structure, or family of objects). For structured data, GridBatch uses splitting the data to segment the data into smaller pieces, similar to database splitting. The GridBatch file system stores files in a fixed number of blocks, each block having a block id (CID). A programmer can access any block, regardless of other blocks in the file system.
In an implementation, the programmer can specify the number of blocks that the GridBatch can designate. In another implementation, a GridBatch administrator specifies the number of blocks that Grid20 Batch can designate, and / or GridBatch determines the number of blocks that GridBatch can designate based on the number of available nodes and / or-or-tras-eonside. -so-from-system-configuration-In one-implementation, the GridBatch file system sets the highest assignable CID to be much greater than N, the number of nodes in the cluster. GridBatch em25 nails a system level lookup table to prescribe the CID mapping for N translation. The translation provides support for dynamically changing the cluster size so that when the configuration deactivates the nodes and additional nodes connect the cluster, the GridBatch file system can automatically rebalance storage and workload. In other words, the file system maintains a CID mapping for the data node, and automatically moves the data to different nodes when the CID for the data node mapping changes (for example, when the data nodes connect and / or leave the GridBatch cluster 102).
In an implementation, GridBatch processes two kinds of data sets: vector and indexed vector. Similar to the records in a database table, a vector includes a set of records that GridBatch considers to be independent of each other. The records in a vector can follow the same scheme, and each record can include several fields (similar to the database columns). Unlike a vector, but similar to an indexed database table, each record in an indexed vector also has an associated index. For example, one of the record fields in the indexed vector could be the associated index of the indexed vector and the index can be of any type of data (for example, string or integer).
When using indexed vectors, the programmer defines how the data should be divided by the blocks using a division function. When a new data record needs to be written, the file system calls the division function to determine the block id and appends the new data record to the end of the block corresponding to the block id. In an implementation, the user-defined division function takes the form: int [] Func.division (index X) where X represents the index for the record to be written and int [] indicates a series of integers. The division function applies a unidirectional function to convert the index into one or more integers in the 1-to-GID range that -indicates the id (s) -of the designated block where the data record must be stored. In another implementation, the division function can take the form: int [] Func.division (distribution key X) where X represents the distribution key indicator for the record to be written to indicate a preferred processor and / or set of processors for use. When using vectors, the GridBatch file system can write each new record to a randomly chosen block.
In an implementation, when a user requests a new file for a new indexed vector to be created, the user provides the file system administrator with a unidirectional user-defined function, which takes the form of int [] Uni-directional Fun. X distribution). The unidirectional function accepts a distribution key as input, and produces one or more integers in the range from 1 to CID. When the new record is written, the file system administrator invokes the unidirectional function to determine which division to write the new record. As a result, GridBatch splits the index vector as new records are processed by the file system administrator.
The work scheduling system includes a master node and multiple slave nodes. The master node can use the logic of the master node to implement the functionality of the master node. A slave node manages the execution of a task assigned to the slave node by the master node. The master node can use the logic of the master node to break a job (for example, a computation) into many smaller tasks as expressed in a program by a programmer. In an implementation, the master node logic distributes tasks across slave nodes in the cluster, and monitors tasks to secure all tasks successfully completed. In an implementation, GridBatch designates data nodes as slave nodes. In this way, when the master node schedules a task, the master node can schedule the task on the node that also maintains the data block to be processed. GridBatch increases computational performance by reducing dependencies on network bandwidth because the GridBatch minimizes data transfers and performs data processing at the data location for the nodes.
GridBatch provides a set of operators called commonly used primitives that the programmer can use to implement computational parallelization. Operators manage the details of job distribution to multiple nodes, so the programmer avoids the burden of addressing the complex issues associated with implementing a parallel programming solution. The programmer introduces a set of operators into a program, in the same pattern as writing a traditional sequential program.
GridBatch provides five operators: distribution, link, convolution, recursion, map. The distribution operator converts a source vector or an indexed vector from source to destination from the indexed vector with a target index. Conversion involves transferring data from a data source node to a data destination node. The distribution operator takes the following form: Vector Distribution (vector V, Func novaFunc.Divisão) where V represents the vector where the data to be converted reside and novaFunc.Divisão represents the division function that indicates the destination data node where GridBatch will generate a new vector. In an implementation, the user-defined division function takes the form int [] novaFunc.Divisão (index X), where X represents the index of the record, and int [] denotes a series of integers. The user-defined division function returns a list of numbers corresponding to the list of destination data nodes. In an implementation, the distribution operator can duplicate a vector on all nodes, so that each node has an exact copy for convenient local processing. Duplication of the vector on all nodes can result when novaFunc. Division returns a list of all data nodes as the target nodes.
The Link operator obtains two indexed vectors and joins the corresponding records where the indexed field matches. GridBatch identifies the corresponding records that have a combined index and involves a user-defined binding function. The user-defined link function can simply join the two records (for example, similar to a database link), but can generally implement any desired function. The link operator takes the form: Link of the Vector (Vector X, Vector Y, Func Func.link) where X and Y represent the indexed vectors to be linked and Func.link represents the user-defined link function to apply the records corresponding in the indexed vectors. The link operator produces a new vector that includes the results of applying the user-defined function. The user-defined link function takes the following form: Link Function Register (Register Z, Register K) where Z and K represent a register of vector X and Y, respectively. When GridBatch invokes the user-defined function, GridBatch can ensure that the indexes for record Z and K match.
GridBatch can perform a distribution operation before executing the link operation so that GridBatch divides the vector X and Y using the division function in the same index field that the Link will subsequently use. The link operator performs the link on each node locally without determining whether GridBatch has distributed or sent data to each node. In an implementation, the liaison operator automatically runs the distribution operator before executing the liaison.
The link operator can be used when there is an exact match in the index field. However, when a programmer wants to identify the inverse result of the bonding operator (for example, identifying mismatch records), every Z record is checked against every K record. The convolution operator identifies Z and K combination records and applies a user-defined function for each combination. The convolution operator provides additional capacity and provides more computational options for the programmer. In an implementation, all computational operations involving two vectors can be performed using the convolution operator. The convolution operator can perform the link function on non-indexed vectors using any field in the vector, even when the link uses an non-indexed field for the link. The convolution operator takes the following form: Vector convolution (vector Xy-vector Y, -func Func.eonv) where X-and Y represent the two input vectors, and Func.conv represents the user-defined convolution function provided by the programmer. The convolution operator produces a new vector as a result. The user-defined function takes the following form: Func.conv register (Z register, K register) where Z and K represent a register of vector X and Y, respectively. The Func.conv function determines whether any action should be taken (for example, it determines whether the Z record matches the K record) and then performs the corresponding action.
GridBatch can execute the Distribution operator before executing the convolution operator so that the GridBatch divides the vector X and Y in the same index field that the convolution can subsequently use15. The convolution operator performs computation on each node locally without determining whether the GridBatch has distributed or sent data to each node. In other implementations, the convolution operator automatically executes the distribution operator before executing the convolution.
As an example, a programmer may wish to determine the number of customers located in close proximity to a retailer's distributors. The GridBatch file system would generate a customer vector that includes a physical location field that indicates the physical location of each customer, and a distribution vector that includes a physical location field that indicates the physical location of each distributor. The programmer can use the GridBatch to join the customer vector and the distributor vector based on the physical location field of both vectors. The programmer can use Func.conv to evaluate the physical distance between each customer and each distributor based on the proximity specified by the programmer, and stores each record satisfying the specified proximity in a result vector.
In an implementation, the recursive functions operator performs a reduced operation, which takes all the records from a vector and joins them into a single result. The actual logical operation performed on the vector records is defined by a user-specified function. In addition, it is an -example-of-reduction operation where all-records of a vector are added together. Classifying another example of the reduction operation where all the records of a vector are checked against each other to produce a desired sequence. The recursive function operator expands the reduction operation through many nodes. Web applications often perform frequent reduction operations (for example, word count, where each word requires a reduction operation to add up the number of appearances), unlike most analytical applications that perform few reduction operations. The reduction operator of most analytical applications becomes a strangler and limits the scalability of an application when a programmer merely needs a classified output for reporting or a few statistics. Many reduction operations exhibit commutative and associative properties, and can be performed in an independent order.
For example, counting the number of occurrences of an event involves the commutative and associative operator known as addition. The order in which the addition occurs does not affect the final result. Similarly, the classification can be independent of order. The GridBatch recursive function operator performs order-independent reduction operations and takes the following form: Recursive Record (Vector X, Func Func.recursive) where X represents the input vector for reduction and Func.recursive represents the recursive function defined by user to apply. The recursive function operator joins the vector into a single record. The user-defined Func.recursive function takes the following form: Func.recursive register (Register Z1, Register Z2) where Z1 and Z2 represent partial results of joins of two sub-parts of the vector X. The Func.recursive function specifies as additionally merge the two partial results.
For example, where the vector X represents a vector of integers and the programmer wants to compute the sum of the integers then the programmer will use the addition function as the user-defined Func.recursive function: Record addition (Register Z1, Register Z2) {return of new Record (value.Z1 () + value.Z2 ())}. The GridBatch will apply the addition function -reeursively on the X-vector records to eventually compute the sum of the integers in the vector.
In another example, vector X includes records that represent classified lists of restrictions and the programmer wants to classify the restrictions for the final report. Table 1 illustrates how GridBatch can implement the user-defined function to classify constraints. The user-defined function joins two classified lists of restrictions into one classified constraint and when the programmer implements the user-defined function to be called again, the user-defined function implements the algorithm to classify the union.
Table 1 - User Defined Function for Classification.
Class.unir Record (Record Z1, Record Z2) {new Record Z;
// next restriction of record Z1 Restriction a = Z1, next ();
// next record restriction Z2 Restriction b = Z2.following (); do {if (a <b) {
Z. attachment (a); a = Z1.next ();
} more {
Z. attachment (b); b = Z2.next ();
}} while (! Z1. empty () &&
! Z2. empty ());
return x;
J
------ The recursive function parallels the reduction operation for many nodes. In addition, the Recursive function minimizes network traffic for operations that require partial results. For example, where a program5 painter needs to identify the top 10 income-producing customers, each node computes the top 10 local customers and sends the results (for example, partial results) to adjacent nodes that in turn combine the partial results with the local result of the receiving node to produce the top 10. Each node only passes the top 10 records to particular adjacent nodes, instead of passing the entire record from each node to a single node performing the reduction operation. In this way, the recursive function operator avoids the demands of large bandwidth and unwanted network traffic, and provides greater computational performance.
The map operator applies a user-defined map function to all records in a vector. The map operator takes the following form: Vector Map (vector V, Func Func.map) where V represents the vector, more specifically the vector records, to which the Func.map will be applied. The user-defined map function can take the following form: Register Func.map (Register X). The user-defined function, Func.mapa, accepts a record from the input vector as an argument and produces a new record for the result vector.
In an implementation, GridBatch tolerates the failures and errors of the slave node by re-executing tasks when slave nodes fail to complete tasks. Each vector block in a vector is duplicated X times in X different slave nodes called backup nodes, where X is a constant that can be specified by the user and / or determined by GridBatch based on configuration, available resources and / or historical observations . During the computation of any operator, if a slave node fails before the slave node completes the assigned task, the master node is informed and the master node starts another process on a slave node that maintains a backup copy of the vector block. The master node identifies a slave node as a failed slave node when the master node does not receive a periodic pulse from the slave node.
------- A-figure - 1— illustrates the system-configuration-GridBatch 100 (GridBatch) that includes a cluster of GridBatch 102, an application 104 and user interface 106. The components of the GridBatch 100 communicate over a network 108 (for example, the Internet, a local area network, a wide area network, or any other network). The GridBatch 102 cluster includes multiple nodes (for example, master node 116 and slave node 120). Each slave node 120 can include a communications interface 113 and memory 118. The GridBatch 100 designates a master node 116, and the remaining nodes, slave nodes (for example, slave node 120). The GridBatch 100 can designate slave nodes as data nodes (for example, data node 134), further described below. Slave node 120 uses slave node logic
160 to manage the execution of slave tasks 158 assigned to slave node 120 by master node 116.
Figure 2 shows an example of Master Node 116. Master node 116 can include a communications interface 211 and memory 215. The
GridBatch 100 uses file system administration logic 222 to manage and store files across all nodes in the GridBatch 102 cluster. In one implementation, file system administration logic 222 segments a large file into small blocks and stores the files. blocks between slave nodes. The system administration logic of file 222 maintains a CID mapping for the data node, and automatically moves the data to different nodes when the CID for data node mapping changes (for example, when the data node becomes joins and / or leaves the GridBatch cluster 102). The GridBatch 100 uses work scheduler logic 230 to coordinate operations between all nodes in the GridBatch 102 cluster.
Among all the nodes in the GridBatch 102 cluster, the GridBatch 100 can designate master node 116 as name node 232, and designates all other nodes to serve as data nodes (for example, data node 134). The name node 232 maintains the file space 238 of the file system 240 240. The name node 232 maintains the mappings of the file vector 242 to the list of corresponding vector blocks, the data nodes designated for eada-bloGOy- and the physical and logical location of each data node. Node name 232 also responds to task requests 244 for the location of a file. In an implementation, node name 232 allocates blocks of large files to data nodes.
Master node 116 fails a task 252 (for example, a computation) as expressed in a program by a programmer in slave tasks (for example, slave task 158) that the work programmer logic 230 distributes among the slave nodes. In one implementation, master node 116 distributes slave tasks across slave nodes in the GridBatch 102 cluster, and monitors slave tasks to secure all successfully completed tasks. In this way, when master node 116 schedules a task 252, master node 116 can program slave tasks (e.g., slave task 158) on the slave node that also maintains the data block to be processed. For example, master node 116 can break task 252 into slave tasks corresponding to slave nodes where the data to be processed reside locally in blocks of the vector, so that GridBatch 100 increases computational performance by reducing bandwidth dependencies minimizing data transfers and performing on-site data processing for nodes.
In one implementation, GridBatch 100 implements master node logic 260 on master node 116 that coordinates communication and interaction between the GridBatch 102 cluster, application 104, and user interface 106. Master node logic 260 can coordinate and control system administrator logic 222 and work scheduler logic 230. The master node logic 260 can hold the GridBatch software library 262 which includes distribution operator logic 264, link operator logic 266, convolution operator logic 268, recursive operator logic 270 and logic from map operator 278. Master node 116 can receive task requests 244 and coordinates the execution of task requests 244 through slave nodes and slave node logic 160.
Figure 3 shows the GridBatch 100 during the processing of a function-call-from-distribution-300 (for example-task request 244) and exercise of the logic of the Distribution Operator 264. In an implementation, the master node 116 receives the call from the distribution function 300 to execute the distribution operator with parameters that include an identifier of the first vector 272 that identifies a first vector to redistribute to obtain blocks of the redistributed vector, redistributed between a set of nodes. For example, the first vector may represent a vector previously distributed with blocks of the distributed vector V1C1 308, V1C2
310, and V1C3 312 between a set of nodes (for example, slave node 1 328, slave node 3 330, and slave node 6 332, respectively). Vector blocks
V1C1 308, V1C2 310, and V1C3 312 include block records of the cor21 vector respondents V1C1R1-V1C1RX 322, V1C2R1-V1C2RY 324 and V1C3R1V1C3RZ 326, respectively.
The master node logic 260 initiates the execution of a division function by generating division tasks 334 in each set of nodes (for example, slave node 1 328, slave node 3 330, and slave node 6 332, respectively) with first vector blocks. Arrow 336 represents a transition to a node state where each node with first vector blocks performs division tasks 334. The records of each block of the vector V1C1 308, V1C2 310 and V1C3 312 of the first vector block can be evaluated by the corresponding division tasks 334 to determine the destination of the vector block designations. For example, each division task 334 can evaluate the block records of the first vector residing in the corresponding slave node to determine a destination of the block location of the vector to redistribute each block record of the first vector. Each division task 334 can create the target of vector block designation files (for example, V1C1F1 338, V1C2F1-V1C2F4-V1C2F3-V1C2F6 340 and V1C3F1V1C3F2- V1C3F5-V1C3F6 342) at the slave location of each corresponding location vector block (for example, destination of the vector block designation) where the block records of the first vector will be redistributed.
Master node 116 can receive task completion notifications for each 334-split-according-task each 334 split task is completed. Master node 116 initiates the execution of a redistribution task by generating redistribution tasks 344 on each slave node (for example, slave node 1 328, slave node 3 330, slave node 4 346, slave node 5 348, node slave 6 332 and slave node 8 350). Arrow 346 represents a transition to a node state where each node that corresponds to the destination of the vector blocks performs redistribution tasks 344. The destination of the vector blocks (for example, V1C1 352, V1C2 354, V1C3 356, V1C4 358 , V1C5 360 and V1C6 362) indicated by the locations of the vector block identified by the vector block designation files (for example, V1C1F1 338, V1C2F1-V1C2F4-V1C2F3-V1C2F6 340 and V1C3F122
V1C3F2- V1C3F5-V1C3F6 342). The redistribution tasks 344 initiate the remote reproduction of the vector block designation files to corresponding target slave nodes to place the vector block designation files on the slave node corresponding to the vector block assigned to the slave node (for example, V1C1F1-V1C3F1-V1C2F1 364, V1C3F2 368, V1C2F3 370, V1C2F4 372, V1C3F5 374, and V1C3F6-V1C3F6 376).
Redistribution tasks 344 initiate a 378 union of records (for example, V1C1R1-V1C1RX 382, V1C2R1-V1C2RY 384, V1C3R1-V1C3RZ 386, V1C4R1-V1C4RQ 388, V1C5R1-V1C5RS 390 and
V1C6R1-V1C6RT 392) located in each vector block designation file corresponding to a particular vector block destination. The arrow 380 represents a transition to a node state in which each node corresponding to the destination of the vector blocks performs a union 378. The union 378 results in the vector blocks redistributed from the first vector, redistributed among the set of nodes. Slave node logic 160 of each slave node sends master communication 116 a completion report indicating the completion status of union 378.
Figure 4 shows the GridBatch 100 when processing a call function call 400 (for example, task request
244) and linker logic exercise 266. In an implementation, master node 116 receives link function call 400 with parameters-that-include-the-identifier-of the first vector-272-and a -identifier of the second vector 274, and a user-defined link function (for example, a user-defined function 276). The identifier of the first vector 272 and an identifier of the second vector 274 identify the first vector and a second vector divided into blocks of the first vector (for example, V2C1 410, V2C2 412 and V2C3 414). The blocks of the first vector and the blocks of the second vector include records of the first vector (for example, V1C1R1V1C1RZ 416, V1C2R8-V1C2RJ 418 and V1C3R4-V1C3RL 420) and records of the second vector (for example, V2C1R3-V2C1R 422, V2C2RK-V2C2RK-V2C2R7-V2C2RK-V2C2R7-V2C2R7-V2 424 and V2C3R4-V2C3RM 426), respectively.
Master node 116 starts the generation of classification tasks (for example, slave tasks 158) locally in the set of nodes (for example, slave node 1 428, slave node 4 430 and slave node 6 432) corresponding to the location of the blocks of the first vector and blocks of the second vector to classify each of the blocks of the first vector and blocks of the second vector to the second vector located in each of the set of nodes. In an implementation, the task of sorting 434 sorts the records of the first vector and the records of the second vector according to an index value of the link index field present in each record of the first vector, of the first vector (for example, V1C1R1IF -V1C1RZIF 438, V1C2R8IF-V1C2RJIF 440 and V1C3R4IF-V1C3RLIF 442) and each record of the second vector, the second vector (for example, V2C1R3IF-V2C1RYIF 444, V2C2R7-V2C2RKIF 446 and V2C3R, respectively) V2C3R The arrow 436 represents a transition to a node state where each node with blocks of the vector performs classification tasks 434.
In one implementation, the classification task 434 compares the index value of the index field present in the records of the first vector and the records of the second vector to determine the records of the first vector and records of the second vector which includes combining index values and applying the user-defined function 276 (for example, a user-defined link function) for records of the first vector and records of the second vector with a combination of the index field values. Classification task 434 performs a combination task 450 that compares the index field values of the index fields of the records of the first vector and records of the second vector. Arrow 452 represents a transition to a node state where each node with blocks in the vector performs 450 matching tasks. Combination task 450 applies user-defined function 276 (for example, a user-defined link function) to records of the first vector and records of the second vector with combination of index field values for corresponding vector blocks (for example , V1C2RBIF 454 and V2C2RPIF 456, and V1C2RBIF 458 and V2C2RPIF 460) to obtain a link function block result (eg, NO JFC1R 462, JFC2R 464 and JFC3R 466). The combination task 450 does not apply the user-defined link function to records of the first vector and records of the second vector when the values of the index field for corresponding vector blocks do not match (for example, V1C1RXIF 468 and V2C1RYIF 470).
The results of the link function block form a result of the link function vector that identifies blocks of the link function vector (for example, JFVC1 476 and JFVC2 478) that includes block records of the link function vector (JFVC1RT 480 and JFVC2R3-JFVC2RN 482) obtained from the connection block results (for example, JFC2R 464 and
JFC3R 466). In an implementation, the slave node logic 160 of each slave node sends master node 116 a completion report that indicates the completion status of classification task 434.
For example, in an implementation, a programmer can use GridBatch 100 to index two vectors, a product vector (for example, first vector identified by the identifier of the first vector 272) indexed by a product id field (for example, V1C1R1IFV1C1RZIF 438, V1C2R8IF-V1C2RJIF 440 and V1C3R4IF-V1C3RLIF 442) and the customer vector (e.g., second vector identified by the identifier of the second vector 274) indexed by the customer id field (e.g. index fields V2C1R3IF-V2C1RYIF 444, V2C2R7-V2C2RKIF 446 and V2C3R4-V2C3RMIF 448). The product vector includes the product id and the -elient-and-corresponding-to-purchased-products-id (for example — index field values). The customer vector maintains the customer id and customer demographic information (for example, values in the index field such as age, address, sex). In the event the programmer wants to know how many people in each age group bought a particular product, the programmer invokes a call function link with the product vector and the customer vector as parameters to obtain a connection result that links the information of the product ID with the customer's demographic information.
In an implementation, in order to ensure greater performance by
GridBatch 100 in processing the link function 400 of the product vector and the customer vector based on the customer id field (for example, index field), the programmer invokes the distribution function call 300 to index the vector of the product by the customer id instead of the product id. The distribution function call ensures that GridBatch 100 distributes the product vector records to nodes in the GridBatch cluster 102 according to the customer id field. GridBatch 100 can then apply user-defined function 276 (for example, a user-defined link function) to each record of the product vector and the customer vector where the values of the customer id field of both the product vector and of the customer vector are equal to obtain the result of the vector of the link function.
Figure 5 shows the GridBatch 100 during the processing of a convolution function call 500 (for example, task request 244) and exercise of convolution operator logic 268. In an implementation, master node 116 receives the function call convolution 500 with parameters that include the first vector 272 identifier and the second vector 274 identifier, and a user-defined convolution function (for example, a user-defined function 276). The identifier of the first vector 272 and an identifier of the second vector 274 identify the first vector and a second vector divided into blocks of the first vector (for example, V1C1 504 and V1C2 506) and the blocks of the second vector <sup>s</sup> 20 (for example, V2C1 508 and V2C2 510) correspond to divided vector blocks distributed through nodes of the GridBatch 102 cluster. The first-vector-blocks and the second-vector-blocks include records of first vector block (for example, V1C1R1-V1C1RZ 512 and V1C3R4-V1C3RL 514) and second vector block records (for example, V2C1R3-V2C1RY 516 and
V2C3R4-V2C3RM 518), respectively.
Master node 116 starts generating convolution tasks (for example, slave tasks 158) locally in the set of nodes (for example, slave node 1 520 and slave node 8 522) corresponding to the location of the first vector blocks and the second vector. The arrow 526 represents a transition to a node state for each node where the master node
116 generates convolution tasks 524. Convolution tasks 524 apply user-defined function 276 (for example, a user-defined conclusion function) locally to the permutations of the block registers of the first vector and the block registers of the second vector (for example, 528 and 530). The user-defined convolution function evaluates each corresponding permutation of first vector block records and second vector block records (for example, 528 and 530) to obtain convolution function evaluation results (for example, 536, 538 , 540 and 542). The arrow 534 represents a transition to a node state for each node where the user-defined convolution function evaluates each permutation of the corresponding vector block records and the second vector block records. The results of evaluating the convolution function can indicate when a permutation of the block records of the first vector and the corresponding second vector block records results in convolution function block result records (for example, CFC1R1-CFC1R3-CFC1R4- CFC1RZ 536 and CFC2R3-CFC2RK 540). The results of evaluating the convolution function can indicate when a permutation of the block records of the first vector and the corresponding second vector block records results in no convolution function block result records (for example, NO CFC1RX 538 and NO CFC2RY 542). The user-defined convolution function can transform the results of the convolution function into block result records of the convolution function (for example, CFVC1R1-CFVC1R3-GFVG1R4-GFVG1 RZ 548-e CFVC2R3-CFVC2RK 550) to obtain function results convolution for each node (for example, slave node 1 520 and slave node 8 522).
For example, in an implementation, a programmer may invoke the call to convolution function 500 to determine the number of customers located in close proximity to a retailer's distributors. The system administration logic for file 222 can include a customer vector (for example, first vector identified by the identifier of the first vector 272) that includes a physical location field that indicates the physical location of each customer and a distributor vector ( for example, second vector identified by the identifier of the second vector 274) which includes a physical location field that indicates the physical location of each distributor. The programmer can invoke the convolution function call 500 to apply a user-defined convolution function (for example, user-defined function 276) to the customer vector and distributor vector based on the physical location field to evaluate the distance between each customer and each distributor and obtain a vector of results of the convolution function. In an implementation, the user-defined convolution function can be expressed as Func.conv. Prior to the convolution call, the customer vector can be divided into blocks of the customer vector (for example, blocks of the first vector - V1C1 504 and V1C2 506) divided through the GridBatch 102 cluster nodes according to the location field physical (for example, index field) present in each of the customer's vector records. The blocks of the distributor vector (for example, blocks of the second vector - V2C1 508 and V2C2 510) can be copied to all nodes in the cluster. This can be achieved by supplying a division function that always returns a list of all nodes to the distribution operator. The user-defined convolution function evaluates the permutations of customer vector records and the distributor vector records residing in corresponding slave nodes, to obtain convolution function block result records. In other words, where the customer vector block has the number Z of records and the
-distributor vector-block has the K-number of registers-the user-defined convolution function can evaluate the number of permutations Z χ K where for each register 1 to Z of the client vector block GridBatch 100 applies the user-defined convolution function to every register 1 to K of the distributor vector block. The result of calling the convolution function performed by each cluster slave node of GridBatch 102 results in blocks of the corresponding convolution function vector to obtain results of the convolution function (for example, slave node 1 520 and slave node 8 522).
Figure 6 illustrates GridBatch 100 when processing a recursive function call 600 (for example, task request
244) and recursive operator logic exercise 270. In an implementation, master node 116 receives the recursive function call 600 with parameters that include the first vector identifier 272 and a user-defined recursive function (for example, function defined by user 276). The first vector identifier 272 identifies the first vector divided into blocks of the first vector (for example, V1C1 604, V1C2 606 and V1C3 610) corresponding to the blocks of the divided vector distributed through the nodes of the GridBatch cluster 102. The blocks of the first vector include block records of the first vector (for example, V1C1R1-V1C1RX 616, V1C1R3-V1C1RJ 618, V1C2R1-V1C2RY 620, V1C2RK-V1C2RN 622, V1C3R4-V1C3RZ 624 and V1C3RG-V1).
Master node 116 starts generating recursive tasks 634 (for example, slave tasks 158) locally in the node set (for example, slave node 1 628, slave node 4 630 and slave node 6 632) corresponding to the location of the blocks of the first vector. Arrow 636 represents a transition to a node state where each block node of the first vector performs the recursive function tasks 634. Recursive function tasks 634 initially apply the user-defined recursive function to the block records of the first vector to produce block results from the immediate recursive vector for each block of the first vector (for example, IRV1C1R1 638, IRV1C1R2 640, IRV1C2R1 642, IRV1C2R2 644, IRV1C3R1 -646 and IR-V-1G3R2-648) -The-recursive-tasks-invoice the user-defined-Feeur-siva function in the immediate recursive vector block results to produce immediate recursive slave node results (for example, IRSN1R 650, IRSN4R 652 and IRSN6R 654).
Recursive tasks communicate a subset of the results of the intermediate recursive slave node (for example, IRSN1R 650) to a subset of the node set (for example, node 4 630) and the invocation iterates over the recursive tasks of the user-defined recursive function in the results intermediates (for example, IRSN1R 650 and IRSN4R 652) to produce progressively less intermediate slave node results (for example, IFIRSN4R 660). Recursive tasks communicate a progressively less subset of intermediate results (for example, IFIRSN4R 660) to a progressively smaller subset of the node set (for example, slave node 6 632) until GridBatch 100 obtains a final recursive result (for example, FRR 668 ) on a final node in the node set.
In an implementation, a subset of the intermediate results communicated by the recursive tasks to a subset of the node set includes half of the intermediate results that produce a progressively smaller subset of intermediate results. Similarly, each subset of progressively less intermediate results subsequently communicated by recursive tasks to a subset of the node set includes one half of the progressively less intermediate results. In an implementation, the recursive operator logic 270 uses network topology information to improve the computational performance of the recursive operator by identifying nearby neighboring slave nodes where intermediate results can be sent and / or retrieved in order to reduce the consumption of the network width. network bandwidth. The programmer, the user and / or 10 can define the factors that determine whether a slave node constitutes a neighbor slave node close to another slave node. Factors that can be used to determine whether a slave node is designated a neighboring slave node can include data transmission times between slave nodes, the number of -hops on the network (for example, -number-of-routers-da -network) between slave nodes, or a combination of data transmission times and network hops.
Figure 6 illustrates how the GridBatch recursive operator logic 270 distributes intermediate results between the slave nodes of the GridBatch 102 cluster. The slave nodes can compute a local intermediate recursive result (for example, IRSN1R 650, IRSN4R 652 and IRSN6R 654). A subset of the slave nodes (for example, slave node 1 628) can transmit the local intermediate recursive result (for example, IRSN1R 650) to a subset of the slave nodes (for example, slave node 4 630). Slave nodes that receive recursive intermediate results daily from other slave nodes can iteratively apply the transmitted intermediate results (for example, IRSN1R 650) with the local intermediate results (for example, IRSN4R 652). Iteratively, even a single slave node (for example, slave node 6 632) produces the final recursive result (for example, FRR 668), a subset (for example, half) of the slave nodes transmits intermediate results to the other half of nodes with local intermediate results (for example, intermediate split results transmitted in local intermediate results). In an implementation, the master node determines the scheme for passing intermediate results to slave nodes in the node set and the number of deployment iterations required to produce a final recursive result (for example, FRR 668).
Figure 7 illustrates the logic flow that GridBatch 100 can obtain to execute the distribution operator. In an implementation, master node 116 receives the distribution function call 300 to execute the distribution operator. In an implementation, the distribution function call 300 can be expressed as Distribution (vector V, func. NovaFunc. Division) The vector V represents the source vector and novaFunc.Division represents a function that determines the location of new nodes for data in the vector V. Figure 7 and the discussion here uses the vector U as a notational aid to explain the redistribution of data. data in vector V. The vector V -contains the same data as the vector U. The distribution-function call 300 results in a remaining vector, possibly divided into new blocks that can be distributed to a different set of we. The master node logic 260 generates a slave task (for example, slave task 158) corresponding to each block of the vector vector V (702). In an implementation, the number of slave tasks is equal to the number of vector blocks in vector V. The slave tasks reside in slave nodes where the corresponding vector blocks (704) reside. Locating slave jobs for slave nodes where the corresponding vector blocks reside minimizes data transfer and avoids scaling issues in the network bandwidth. The slave nodes invoke the slave node logic 212 to generate output files that correspond to blocks of the vector vector U where GridBatch 100 will redistribute records of the vector V (706). Slave node logic 160 evaluates each record block of the corresponding vector of V to determine the block identifier of vector U where GridBatch 100 will redistribute the record. The slave node logic 160 writes the record to the output file that corresponds to the vector vector block U where GridBatch 100 will redistribute the vector V record.
As each slave task completes the evaluation of the records of the corresponding vector blocks of V, each slave task notifies the master node logic 260 of the completion status of the slave task and the location of the output files corresponding to the blocks of the vector vector U ( 708). The master node logic 260 generates new slave tasks on slave nodes where GridBatch 100 will redistribute vector vector blocks V to vector vector blocks U (710). Each slave task receives a list of output file locations that include blocks of the U vector that correspond to the slave node corresponding to the slave task and retrieves the output files to the slave node (for example, using a remote copy operation, or another file transfer). Each slave task joins the output files into corresponding U vector blocks and notifies master node logic 260 of the completion status of the slave task (712). In one implementation, the distribution function call 300 distributes all-of-the-first-register-vectors to all available slave nodes- For example, the newFunc.Division of the distribution function call 300 expressed as Distribution (vector V, func. novaFunc.Divisão) can direct GridBatch 100 to distribute each record of vector V to all available slave nodes to duplicate vector V in all available slave nodes.
Figure 8 shows the logic flow that GridBatch 100 can obtain to execute the connection operator. In an implementation, the master node logic 260 receives the link function call 400 to link the vector X and the vector Y. In an implementation, the link function call 400 can be expressed as the link of the Vector (vector X, vector Y, Func.
Connection) (802). The master node logic 260 generates a slave task that corresponds to a vector block number (for example, vector block id), where the 222 file system administration logic divides vector X and vector Y into one equal number of vector blocks and 222 file system administration logic designates X vector blocks and Y vector blocks with corresponding block numbers or vector block ids (804). For example, the system administration logic of file 222 can assign a particular block id to either a vector X block or a vector Y block residing in a corresponding slave node. In an implementation, the slave task classifies, according to an indexed field value, the X vector block records and the Y vector block records residing in the corresponding slave node (806). The slave task invokes the slave node logic 160 and evaluates the indexed field value of the X vector block records and the Y vector block records. If the indexed field values of the X vector block records and records of the Y vector block are equal (808), GridBatch 100 invokes a user-defined link function (for example, user-defined function 276). In one implementation, the user-defined link function can be expressed as Link Function Register (Z Register, K Register) that links the X vector block records and the Y vector block records (814). If the slave node logic 160 evaluates the value of the indexed-field of record Z-of the block of the vector-X-be less than -the value of the indexed field of record K of the block of the vector of Y then the slave node logic 160 evaluates the next Z record of the X vector block with the value of the indexed K record field of the Y vector block (810). If the slave node logic 160 evaluates the value of the indexed field of record Z of the block of the vector X is greater than the value of the indexed field of record K of the block of the vector of Y then the slave node logic 160 evaluates the next record K of the Y vector block with the value of the Z record indexed field of the X vector block (812). Slave node logic 160 evaluates every Z record in the X vector block and K record in the Y vector block (816).
Figure 9 shows the logic flow that GridBatch 100 can obtain to execute the convolution operator. In one implementation, the master node logic 260 receives the call from the convolution function 500 to process the vector X and vector Y (902). In an implementation, the convolution function call 500 can be expressed as Vector convolution (vector X, Vector Y, Func Func.conv), where Func.conv is the user-specified convolution function. For each register 1 to Z of the vector blocks of vector X, the master node logic 260 applies a user-defined convolution function (for example, user-defined function 276), expressed as Link Function Register (Register Z, Register K ) for records 1 to K of vector blocks of vector Y (904). In other words, if a vector block of vector X has the number Z of records and a vector block of vector Y has the number K of records, the user-defined convolution function evaluates the number of permutations Z x K of pairs from register. The slave node logic 160 applies the user-defined convolution function for each register 1 to K of the vector Y (906) with every register 1 to Z of the block of the vector X (908).
Figure 10 shows the logic flow that GridBatch 100 can obtain to execute the recursive operator. In an implementation, the master node logic 260 receives the recursive function call 600 for the recursive vector X. In an implementation, the recursive function call 600 can be expressed as Recursive Record (Vector X, Func Func.recursive). The master-clock-2-60 ~ generates Fecursive-operation-slave-tasks corresponding to each vector block residing in corresponding slave nodes (1002). Slave jobs invoke slave node logic 160 to reduce (for example, join) the first record and the second vector block records of vector X residing in corresponding slave nodes. Slave node logic 160 stores the intermediate recursive result (for example, joining) (1004). Slave node logic 160 assesses whether more vector X vector block records exist (1006) and joins the next vector X vector block record to the intermediate join result (1008). Since the slave node logic 160 obtains the result of intermediate union of the vector blocks of the vector X, each slave task notifies the master node logic
260 completion status of the slave task (1010). A subset of slave jobs (for example, half) send intermediate join results for the remaining slave jobs (for example, the other half) with local intermediate results. The subset of slave tasks that receive the intermediate join results join the intermediate join tasks with the local intermediate join results (1012). Slave nodes with intermediate union results iteratively double the intermediate union results in fewer slave nodes, until the slave nodes join the progressively smaller number of intermediate union results in a final union result residing in a slave node (1014).
Figure 11 illustrates the GridBatch 100 when processing a map function call 1100 (for example, task request 244) and exercising map operator logic 278. The map operator can be expressed as Vector Map (vector V, Func Func.map) where V represents the vector, more specifically the vector records, for which Func.map will be applied to obtain a new vector of mapped vector V records. The map operator allows the user to apply a user-defined function to all records in a vector. In an implementation, the master node logic 260 receives the map function call 1100 with parameters that include an identifier of the first vector 272 and -a user-defined map function (for example, a function user-defined 276). The first vector identifier 272 identifies the first vector divided into blocks of the first vector (for example, V1C1 1104, V1C2 1108 and V1C3 1110) corresponding to divided vector blocks distributed through the nodes of the GridBatch cluster 102. The blocks of the first vector include block records of the first vector (for example, V1C1R1 1116, V1C1RX 1118, V1C2R1 1120, V1C2RY 1122, V1C3R4 1124, and V1C3RZ 1126).
The master node 116 initiates a generation of map tasks 1134 (for example, slave tasks 158) locally in the node set (for example, slave node 1 1128, slave node 4 1130 and slave node 6 1132) corresponding to the location of the blocks of the first vector. The arrow 1136 represents a transition from a node state where each node with first vector blocks performs map tasks 1134 (for example, map tasks running in parallel 1150, 1152 and 1154). Map tasks 1134 apply the user-defined map function to each vector block record to produce the mapped vector block records that form mapped vector blocks of the vector Μ. Arrow 1158 represents a transition to a node state where each node with blocks from the first vector includes corresponding mapped vector blocks (for example, VMC1 1160, VMC2 1162, and VMC3 1164) with block records from the corresponding mapped vector (for example, example, VMC1R1 1166, VMC1RX 1168, VMC2R1 1170, VMC2RY 1172, VMC3R4 1174, and VMC3RZ 1176).
For example, a 1180 sales record vector can include a customer ID, product ID, and purchase field date, along with several other fields. However, for a particular analysis, only two fields in the sales record vector may be of interest, such as the customer ID and the product ID. For efficient processing performance, a programmer can invoke the map function call 1100 to execute the map operator to extract just the Customer ID and Product ID fields from the sales record vector; the map function call 1100 can be expressed as follows: Vector Ve-tor.new --- Map (Vector.Register ^ sale, -fender). The user-defined fender-function analyzes each record in the 1180 sales record vector to produce new records that only include customer ID and product ID fields in the Vector.new 1182 records.
Figure 12 shows the logic flow that GridBatch 100 can obtain to execute the map operator. Master node logic 260 receives map function call 1100 for map vector V (1202). The master node logic 260 generates corresponding slave tasks for each vector block of the vector V (1204). Slave jobs invoke slave node logic 160 to locate each vector block in vector V assigned to corresponding slave nodes (1206). For each vector block in vector V, slave node logic 160 applies user-defined Func.map to each vector block record to obtain mapped vector block records that form a mapped vector block of vector M (1208 ). Since slave node logic 160 applied to Func.map to each block record of vector V, each slave task notifies master node logic 260 of the completion status of the slave task and the location of the mapped vector block. corresponding to Μ. The map operator ends successfully when the slave nodes notify the master node that all slave tasks are finished (1210). The mapped vector blocks of vector M combine to form a new vector M.
The additional operators that GridBatch provides produce unexpectedly good results for parallel programming techniques. In particular, each operator provides significant advantages over previous attempts at application parallelization. The unexpected good results include significant flexibility, efficiency and applicability of programming, additional to extraordinarily difficult problems faced by modern businesses, particularly with huge amounts of data that must be processed in a realistic time structure to achieve significant results.
The Redução.mapa programming model implements a unitary programming construction. In particular, a Map function is -always-paired with a Reduce-function. On the other hand, GridBatch provides multiple independent operators: Resource, Convolution, Link, Distribution and Map that a programmer can use in virtually any order or sequence to build a complex application that runs in parallel across many nodes. In addition, the GridBatch support structure implements the user with defined functions specified for the independent operators through which the programmer can confer an immense degree of special functionality. Such user-defined functions include a division function to determine how to break a vector into blocks, a unidirectional function to distribute vector blocks between nodes, a link function to specify how to combine resources, a convolution function to support the operator link, a recursive function that specifies how to join partial results from the recursive operator, and a map function to apply to vector records.
A number of implementations have been described. Nevertheless, it will be understood that various modifications can be made without departing from the spirit and scope of the invention. Accordingly, other implementations are within the scope of the attached claims.
20 members in 6 offices
Priority claims1
| Document | Office | Kind | Date |
|---|---|---|---|
| 90629307 | United States of America | A |
Members20
| Document | Office | Kind | |
|---|---|---|---|
| CA2639853A1 | Canada | A1 | |
| US2009089544A1 | United States of America | A1 | |
| US2009089560A1 | United States of America | A1 | |
| CN101403978A | China | A | |
| EP2045719A2 | European Patent Office (EPO) | A2 | |
| EP2045719A3 | European Patent Office (EPO) | A3 | |
| BRPI0805054A2This record | Brazil | A2 | |
| AR068645A1 | Argentina | A1 | |
| CA2681154A1 | Canada | A1 | |
| EP2184680A2 | European Patent Office (EPO) | A2 | |
| EP2184680A3 | European Patent Office (EPO) | A3 | |
| CN101739281A | China | A | |
| US7917574B2 | United States of America | B2 | |
| US7970872B2 | United States of America | B2 | |
| CN101403978B | China | B | |
| CA2681154C | Canada | C | |
| CA2639853C | Canada | C | |
| CN101739281B | China | B | |
| EP2184680B1 | European Patent Office (EPO) | B1 | |
| BRPI0805054B1 | Brazil | B1 |
11 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Lapse because of non-payment of annual fees (definitively: art 78 iv lpi, resolution 113/2013 art. 12)LapsedEM VIRTUDE DA EXTINCAO PUBLICADA NA RPI 2847 DE 29-07-2025 E CONSIDERANDO AUSENCIA DE MANIFESTACAO DENTRO DOS PRAZOS LEGAIS, INFORMO QUE CABE SER MANTIDA A EXTINCAO DA PATENTE E SEUS CERTIFICADOS, CONFORME O DISPOSTO NO ARTIGO 12, DA RESOLUCAO 113/2013.B24J | B24J | |
| Lapse acc. art. 78, item iv - on non-payment of the annual fees in timeLapsedREFERENTE A 17A ANUIDADE.B21F | B21F | |
| Patent or certificate of addition of invention granted [chapter 16.1 patent gazette]GrantedPRAZO DE VALIDADE: 10 (DEZ) ANOS CONTADOS A PARTIR DE 12/11/2019, OBSERVADAS AS CONDICOES LEGAIS. (CO) 10 (DEZ) ANOS CONTADOS A PARTIR DE 12/11/2019, OBSERVADAS AS CONDICOES LEGAISB16A | B16A | |
| Decision: intention to grant [chapter 9.1 patent gazette]B09A | B09A | |
| Others concerning applications: alteration of classificationAS CLASSIFICACOES ANTERIORES ERAM: G06F 15/80 , G06F 19/00 , G06F 9/28B15K | B15K | |
| Patent application procedure suspended [chapter 6.1 patent gazette]B06A | B06A | |
| Objections, documents and/or translations needed after an examination request according [chapter 6.6 patent gazette]B06F | B06F | |
| Formal requirements before examination [chapter 6.20 patent gazette]B06T | B06T | |
| Requested transfer of rights approvedB25A | B25A | |
| Requested transfer of rights approvedB25A | B25A | |
| Publication of a patent application or of a certificate of addition of invention [chapter 3.1 patent gazette]B03A | B03A |
Numbers
- Application
- 8050546
Titles2
- Portuguese
- infraestrutura para programação paralela de grupos de máquinas
- English
- infrastructure for parallel programming of machine groups
Classification
- CPC, 2
- G06F8/45
- G06F16/24532
- IPC, 3
- G06F15 80
- G06F9 28
- G06F19 00