Accelerating Data Shuffling in MapReduce Framework with Scale-up NUMA Computing Architecture

Xiang Cao, Kewal Keshaorao Panchputre, David Hung-Chang Du · 2016

Nowadays, MapReduce has become very popular in many applications, such as high performance computing. It typically consists of map, shuffle and reduce phases. As an important one among these three phases, data shuffling usually accounts for a large portion of the entire running time of MapReduce jobs. MapReduce was originally designed in scale-out architecture with inexpensive commodity machines. However, in recent years, scale-up computing architecture for MapReduce jobs has been developed. Some studies indicate that in certain cases, a powerful scale-up machine can outperform a scale-out cluster with multiple machines. With multi-processor, multi-core design connected via NUMAlink and large shared memories, NUMA architecture provides a powerful scale-up computing capability. Compared with Ethernet connection and TCP/IP network, NUMAlink has a much faster data transfer speed which can greatly expedite the data shuffling of MapReduce jobs.

Read the paper · More papers on PaperTik