HPAT – A compiler-based framework for big data in Python
github.com
github.com
> Data frames in scripting languages are essential abstractions for processing structured data. However, existing data frame solutions are either not distributed (e.g., Pandas in Python) and therefore have limited scalability, or they are not tightly integrated with array computations (e.g., Spark SQL). This paper proposes a novel compiler-based approach where we integrate data frames into the High Performance Analytics Toolkit (HPAT) to build HiFrames. It provides expressive and flexible data frame APIs which are tightly integrated with array operations. HiFrames then automatically parallelizes and compiles relational operations along with other array computations in end-to-end data analytics programs, and generates efficient MPI/C++ code. We demonstrate that HiFrames is significantly faster than alternatives such as Spark SQL on clusters, without forcing the programmer to switch to embedded SQL for part of the program. HiFrames is 3.6x to 70x faster than Spark SQL for basic relational operations, and can be up to 20,000x faster for advanced analytics operations, such as weighted moving averages (WMA), that the map-reduce paradigm cannot handle effectively. HiFrames is also 5x faster than Spark SQL for TPCx-BB Q26 on 64 nodes of Cori supercomputer.
for i in range(num_iterations):
some_function(i)
A library like Spark can execute the function in parallel, but must return control to the host language for each iteration. That leads to a restart of the distributed environment.HPAT is a compiler. The function is executed in parallel, just like in Spark, but the distributed environment doesn't return control the host language between iterations. Instead, the loop is run in its entirety in the distributed environment.
Spark's driver must farm work to executors during each iteration. HPAT removes that runtime overhead.
[0] https://www.dursi.ca/post/hpc-is-dying-and-mpi-is-killing-it...
"Similarly, function calls should also be deterministic. The below example is not supported since function f is not known in advance"
What if the if statement is valuable at compile time?
What kind of existing infrastructure would you need to actually run this? I assume you would need to install agents that understand MPI on all your nodes and run some sort of centralised master agent to send out messages and supervise the whole thing. Is there a standard way to do this?
A vanilla compiler will suffice for MPI code, so you don't need anything special at build time.