Complex Query Processing with MapReduce in a Multi-Terabyte Logging Cluster

Conrado Plano · Repository for Publications and Research Data (ETH Zurich) · 2011

Hadoop, a MapReduce open source implementation, is the de facto tool for big data analysis in distributed environments.The logging cluster presented here processes around 500 terabytes of logs per week, which are stored in a cluster of logging servers and need to be processed using complex queries.This paper introduces modifications and new tools to Hadoop, that enhance it and prepare it for the required environment and the related use cases.A new input format to read such logs is defined and HLDFS is introduced.HLDFS is a new Hadoop Lightweight Distributed File System model and implementation that allows the use of already distributed files as input in Hadoop.The notion of Shared Scans in Hadoop is also modeled and implemented in order to counteract the existing bottleneck in the map phase execution.The system is successfully tested under production conditions and delivers promising results on the input format and the HLDFS file system.

Read the paper · More papers on PaperTik