scieee AI-readable full text Open interactive document viewer

Repositorio Institucional de Documentos

Abstract

Cloud computing has opened doors to a new era of enterprises that harness the new Cloud enabled business. More and more novel applications are leveraging this paradigm every day, which translates to a never seen increase in the amount of stored data. This phenomenon is commonly known as Big Data; the presence of rapidly expanding high-volume data sets. Many of these applications bring new challenges to databases and therefore the scalability of the Cloud-based databases has become a top-research issue of the Cloud Computing infrastructure. As an alternative to the well-known relational databases, NoSQL databases have born to fit Big Data application requirements. Traditional relational databases as they are often implemented are not sufficient anymore for Internet scale distributed systems dealing with Big Data. Nevertheless, NoSQL have proved to be robust in Big Data applications. The purpose of this project is to scaling out the data of a San Diego company that produces software for wireless multimedia, which is currently implemented on a MySQL cluster. In order to improve the performance of the computation, we propose a solution using Apache HBase, a NoSQL database. This final project proposes implementations as well as comparison details of a number of computation techniques conducted in HBase along with different open-source distributed computing components such as Hadoop HDFS and MapReduce, and presents benchmarks of our developed solution. Cerdán Lázaro, Mario; Heljanko, Keijo

Full text

Repositorio de la Universidad de Zaragoza – Zaguan http://zaguan.unizar.es Anexo PFC Scalability of a cloud-based data store: Improving HBase performance Autor Mario Cerdán Lázaro Director Keijo Heljanko Ponente José Ángel Bañares Bañares INGENIERÍA INFORMÁTICA 2014 2 / 2 Repositorio de la Universidad de Zaragoza – Zaguan http://zaguan.unizar.es 1. Anexo A continuación se presenta la memoria original del proyecto. La memoria es autoexplicativa y contiene: - Titulo - Resumen del proyecto (Abstract). - Lista de abreviaciones y acrónimos. - Lista de contenidos. - Lista de figuras. - Memoria - Bibliografía Este anexo ha sido realizado usando la herramienta Latex y siguiendo el formato de la universidad de Helsinki (Aalto University), por lo que este varía con respecto a la versión de EINA (Universidad de Zaragoza). Aalto University School of Science Degree Programme in Computer Science and Engineering Mario Cerdan Scalability of a Cloud-Based Data Store: Improving HBase performance Final Project Espoo, November 13, 2013 Supervisor: Assoc. Prof. Keijo Heljanko Advisor: Assoc. Prof. Keijo Heljanko Aalto University School of Science Degree Programme in Computer Science and Engineering ABSTRACT OF FINAL PROJECT Author: Mario Cerdan Title: Scalability of a Cloud-Based Data Store: Improving HBase performance Date: November 13, 2013 Pages: 94 Major: Computer Science and Engineering Code: T-110 Supervisor: Assoc. Prof. Keijo Heljanko Advisor: Assoc. Prof. Keijo Heljanko Cloud computing has opened doors to a new era of enterprises that harness the new Cloud enabled business. More and more novel applications are leveraging the Cloud Computing paradigm every day, which translates to a never seen increase in the amount of stored data. This phenomenon is commonly known as Big Data; the presence of rapidly expanding high-volume data sets. Many of these applications bring new challenges to databases and therefore the scalability of the Cloud-based databases has become a top-research issue of the Cloud Computing infrastructure. As an alternative to the well-known relational databases, NoSQL databases have born to fit Big Data application requirements. Traditional relational databases as they are often implemented are not sufficient anymore for Internet scale distributed systems dealing with Big Data. Nevertheless, NoSQL databases have proved to be robust in Big Data applications. The purpose of this project is to scaling out the data of a company that produces software for wireless multimedia, which is currently implemented on MySQL cluster. In order to improve the performance of the computation, we proposes a solution using Apache HBase, a NoSQL database. This final project proposes implementations as well as comparison details of a number of computation techniques conducted in HBase along with different opensource distributed computing components like Hadoop HDFS and MapReduce, and presents benchmarks of our developed solution. Keywords: Cloud-based, datastore, NoSQL, HBase, Hadoop, MapReduce, CAP, Skew Data, YCSB Language: English 2 Acknowledgements This work would not have been completed without help and support of many individuals. I would like to thank to my supervisor Keijo Heljanko for providing me an opportunity to conduct my Thesis under his invaluable guidance and support over the course of it. I am grateful to Aalto University for giving me the chance of finishing my studies in Finland and for the opportunity of using Triton cluster for this Thesis. I would like also to thank to my roommates at Aalto: Bailo, Canellas and Guillermo for all the unforgettable moments we shared. This Thesis is dedicated to the three pillars of my life: my mother and her indefatigable support, my father and his efforts of making me happy no matter what happens and to my brother Jorge because without him I would be lost. To my family, to whom I owe my life. Thanks. Espoo, November 13, 2013 Mario Cerdan 3 Abbreviations and Acronyms SaaS Software as a Service PaaS Platform as a Service IaaS Infrastructure as a Service AWS Amazon Web Services GCP Google Cloud Platform GAE Google App Engine DBMS Database Management System RDBMS Relational Database Management System GFS Google File System HDFS Hadoop Distributed File System CAP Consistency, Availability and Partition Tolerance BASE Basically, Available, Soft state, Eventually consistent WAL Write-Ahead Log YCSB Yahoo! Cloud Serving Benchmark XML eXtensible Markup Language 4 Contents Abbreviations and Acronyms 4 1 Introduction 10 1.1 Cloudproviders.......................... 13 1.1.1 Amazon .......................... 13 1.1.2 Google........................... 15 1.1.3 Microsoft ......................... 16 1.1.4 RackSpace......................... 17 1.2 Behind the big Cloud providers . . . . . . . . . . . . . . . . . 17 1.2.1 IaaS ............................ 18 1.2.2 PaaS............................ 18 1.2.3 SaaS............................ 19 1.3 BigData.............................. 20 2 Background 23 2.1 Datastores: From SQL to NoSQL systems . . . . . . . . . . . 23 2.1.1 The basic principles of NoSQL . . . . . . . . . . . . . . 25 2.1.2 Key features of NoSQL Datastores . . . . . . . . . . . 28 2.1.3 Types of NoSQL Datastores . . . . . . . . . . . . . . . 28 3 Technical background 30 3.1 Column-Oriented Datastores . . . . . . . . . . . . . . . . . . . 30 3.2 HBase ............................... 30 3.2.1 DataModel........................ 31 3.2.2 Storage .......................... 32 3.2.3 Architecture........................ 32 3.2.3.1 Storage layer . . . . . . . . . . . . . . . . . . 33 3.2.3.2 Server layer . . . . . . . . . . . . . . . . . . . 34 3.2.3.3 Client layer . . . . . . . . . . . . . . . . . . . 34 3.2.4 Write............................ 35 3.2.5 Read............................ 36 5 3.2.6 Delete ........................... 36 3.2.7 HBaseAPI ........................ 36 3.2.8 HBase properties . . . . . . . . . . . . . . . . . . . . . 40 3.3 HDFS ............................... 41 3.4 MapReduce ............................ 41 3.4.1 Hadoop MapReduce . . . . . . . . . . . . . . . . . . . 42 4 Environment 44 4.1 Triton ............................... 44 4.2 Cloudera’s Distribution Including Apache Hadoop - CDH . . . 45 4.3 MySQL .............................. 45 4.4 Yahoo! Cloud Serving Benchmark - YCSB . . . . . . . . . . . 45 4.5 Introducing our Dataset . . . . . . . . . . . . . . . . . . . . . 45 4.5.1 HBase storage schema design . . . . . . . . . . . . . . 46 5 Design and Implementation - Import 48 5.1 HBase cluster at a glance . . . . . . . . . . . . . . . . . . . . . 49 5.2 HBase: Tuning parameters for a write-heavy cluster . . . . . . 49 5.2.1 Hadoop baseline . . . . . . . . . . . . . . . . . . . . . . 52 5.3 The import experiment . . . . . . . . . . . . . . . . . . . . . . 53 5.3.1 First approach: An HBase client . . . . . . . . . . . . . 53 5.3.2 Second approach: A multithread HBase client . . . . . 56 5.3.3 Third approach: Using the MapReduce algorithm . . . 58 5.3.3.1 Building the solution . . . . . . . . . . . . . . 58 5.3.3.2 Third approach second version: Compression 60 5.3.3.3 Third approach third version: Pre-creating regions...................... 62 5.3.4 Fourth approach: Coping with skewed data . . . . . . 64 5.4 Performance Tuning Hadoop . . . . . . . . . . . . . . . . . . . 66 6 Design and Implementation - Retrieval 71 6.1 Random reads in HBase . . . . . . . . . . . . . . . . . . . . . 71 6.1.1 Random reads in our heavy-write cluster . . . . . . . . 71 6.1.2 Studying random read performance . . . . . . . . . . . 72 6.1.3 Proceding with random read . . . . . . . . . . . . . . . 73 7 Design and Implementation - Benchmarking 77 7.1 Benchmarking: HBase vs MySQL . . . . . . . . . . . . . . . . 77 7.1.1 Loadphase ........................ 78 7.1.2 Reads mixed with Updates . . . . . . . . . . . . . . . . 79 7.1.3 Predominant Reads . . . . . . . . . . . . . . . . . . . . 80 6 8 Conclusions 83 8.1 Futurework............................ 83 8.2 Discussion............................. 84 7 CHAPTER 1. INTRODUCTION 14 Amazon EC2 provides users with total control of their computing resources. (b) Amazon Elastic MapReduce (EMR) allows businesses, researchers, data analysts, and developers to easily and cheaply process vast amounts of data. It uses a hosted Hadoop framework running on the web-scale infrastructure of EC2 and Amazon S3. 2. Storage & Content Delivery: (a) Simple Storage Service (S3) provides Web Service based storage. (b) Amazon Glacier, Provides a very low cost long-term storage option (when compared to its S3 service). High redundancy and availability, but very long access latency. Ideal for archiving data. (c) AWS Storage Gateway is an iSCSI block storage virtual appliance with cloud-based backup. (d) Amazon Elastic Block Store (EBS) provides persistent block-level storage volumes for EC2. 3. Database: (a) Amazon DynamoDB provides a scalable, low-latency NoSQL online Database Service backed by Solid State Disks (SSDs). (b) Amazon ElastiCache provides in-memory caching for web applications. This is Amazon’s implementation of Memcached. (c) Amazon Relational Database Service (RDS) provides a scalable database server with MySQL, Oracle, and SQL Server support. (d) Amazon Redshift provides petabyte-scale data warehousing with column-based storage and multi-node compute. (e) Amazon SimpleDB, allows developers to run queries on structured data. It operates in concert with EC2 and S3 to provide ”the core functionality of a database.” 4. Deployment (a) AWS Elastic Beanstalk provides quick and easy deployment and management of applications in AWS cloud. It is the main Amazon’s PaaS solution. Elastic Beanstalk harnesses AWS services to complete its task successfully. Elastic Beanstalk takes care of the application’s deployment: capacity provisioning, load balancing request, automatic scaling, etc. It supports Node.js, PHP, Python, Ruby, .NET and Java. CHAPTER 1. INTRODUCTION 15 Amazon Web Services is used by many cloud companies to provide new cloud services, including RightScale providing IaaS, Heroku providing PaaS or Dropbox providing SaaS. 1.1.2 Google Google infrastructure has been built and continues being built to work on datacenters of commodity hardware as opposed to high-end hardware due to what is called the economy of scale: costs are reduced and overall computing power is maximized. Google has managed to develop a fault-tolerant elastically scalable system to work on several datacenters of commodity hardware, letting them to offer cheap Cloud Services. With that premise, Google entered the Cloud Computing business in 2008 with the Google App Engine offering. Whilst cloud provisioning is not their core business, they have given to the world numerous and important contributions in this field. Google App Engine (GAE) is Google’s PaaS solution. GAE allows users to host and run web applications and store data in Google-managed datacenters distributed around the world. It supports Java, Python, PHP and Google’s Go programming languages. Users develop their applications on their local machines before uploading the applications to GAE. Then, it is GAE who takes care of the provisioning and deployment of the application uploaded on the infrastructure. Also automatic scaling is done during the life of the deployed application. Users do not need to keep an eye on the servers as this is GAE infrastructure’s job. GAE software development toolkit (SDK) is provided by Google in order to allow users to develop their solutions in a simulated GAE’s environment. Google also supplies APIs that can be used to integrate Google services with developer applications. On June 2012, Google announced Google Compute Engine (GCE) [cite] at Google IO. GCE is Google’s IaaS solution. GCE allows users to run large-scale CPU works on Linux virtual machines hosted on Google’s infrastructure. GCE provides all resources through the Google APIs Console, a collection of Google’s APIs. GCE is very similar to Amazon’s EC2 solution, both provides scalable CPU capacity. After that, Google unified its cloud solutions under the name of Google Cloud Platform (GCP) [37]. Google Apps [36] is the Google’s SaaS implementation launched on 2006 and entirely based on Google’s own infrastructure. It is a set of several Web applications which offer an online alternative to traditional offline office software. But Google Apps is not just an online office suite, but also a solution that allows users to communicate and collaborate between their projects easily. With Google Mail users can communicate with emails, online messaging CHAPTER 1. INTRODUCTION 16 and voice or video calls. Thanks to Google Docs users can create, edit, delete their documents, spreadsheets or slides. Google Calendar is a powerful online calendar application. Google Web Pages allows for publication Web pages. Google Drive offers file storage and synchronization to users. Recently, in the Google IO 2013 meeting, Google added to their Google Apps collection Google Hangouts. It creates online video meetings with a click, allowing users to work with clients or partners in real time. The main advantage of Google Apps is that everything is online, users do not need to install any software locally, they just need a computer with Internet connection and a Web browser to interact with Google Apps. 1.1.3 Microsoft When analyzing this field, Microsoft has always something to offer. Windows Azure [59] is Microsoft’s platform for Cloud Computing. It was announced at the Professional Developers Conference in October 2008, and commercially available since February 2010. In the beginning, Windows Azure supported only .NET development. However, now it supports many different programming languages, tools and systems, not only Microsoft-specific ones, but also third-party tools, such as Java, Ruby, PHP, Node.js and C++. Windows Azure is hosted in Microsoft-managed datacenters distributed around the world. Microsft states that Windows Azure enables users to deploy their applications within minutes and scale in/out them to any size in a fully automated way. Windows Azure provides an API built on REST, HTTP, SOAP and XML that allows users to communicate with Windows Azure services easily. Microsoft also provides open-source Windows Azure client libraries for multiple programming languages. These SDKs help users build, deploy and manage their Windows Azure applications. Nowadays, Windows Azure Platform offers Infrastructure as a Service (IaaS) features that complement their initial offering of PaaS features. The Azure platform provides three distinct computing service models: 1. Windows Azure Web Sites, which is their PaaS service for Web hosting. Users can create web sites in PHP, .NET, Python and Node.js and deploy them using git, FTP or TFS. Web Sites supports horizontal scalability of web sites, from shared single instances to dedicated large instances. 2. Windows Azure Cloud Services, the traditional Microsoft PaaS service offering. Cloud Services are containers of hosted applications. Applications execute in virtual machines, also called instances, running Win- CHAPTER 1. INTRODUCTION 17 dows Server OS. Windows Azure itself manages the instances. Cloud Services allows to create scalable and reliable applications. 3. Windows Azure Virtual Machines, which comprises the Microsoft IaaS solution for their public cloud. Users create Virtual machines on demand. Unlike Windows Azure Cloud Services, users have total control of their created virtual machines. The Virtual machines offering includes Windows Server images as well as Linux distributions images provided by Microsoft partners. Windows cloud offerings are not just PaaS and IaaS solutions, but also SaaS tools. Windows Live is a Windows’s SaaS which integrates search, email and a social network system. Skydrive is the cloud hosting Microsoft SaaS model. And Microsoft SharePoint is Microsoft collaborative cloud system, which allows multiple users to work together in real time. 1.1.4 RackSpace Behind these three big players, RackSpace [65] is one of the strongest cloud providers companies. Powered by an offering mostly based on the OpenStack Cloud Computing platform, Rackspace is a mostly open source based alternative to Amazon, Google and Windows Azure. One of its strengths is its ability to roll out the latest Openstack features, thus continuosly improving its functionality. Rackspace offers three different Cloud Services: Cloud Servers, their IaaS solution, Cloud Sites, which is their PaaS solution, and Cloud Files for storing files in the cloud, which is their SaaS solution. As Eric Savitz states in Forbes article 1, in the last quarter of 2012, Rackspace total server amount has increased from 89.051 to 90.524, along with an increase of their total customers from 197.635 to 205.538. 1.2 Behind the big Cloud providers Nowadays Amazon, Google, Microsoft are the biggest Cloud providers over the world, but this does not mean there are no other competitors offering excelent products in the Cloud field. Here we show the most promising cloud providers categorized by their offerings of resources: IaaS, PaaS and SaaS. 1Rackspace Slides After Hours As Q4 Revs Miss Street Views -http://www.forbes.com/sites/ericsavitz/2013/02/12/ rackspace-slides-after-hours-as-q4-rev-miss-street-views/ CHAPTER 1. INTRODUCTION 18 1.2.1 IaaS Infrastructure as a Service cloud segment accounted for the majority of total market revenue in 2012 with more than half of the total public cloud market share [98]. The main responsibles for this impressive market quota are Google with its Google Cloud Platform, Amazon and Microsoft Azure, all of them deeply studied in the previous section 1.1. 1.2.2 PaaS Among the subcategories of Cloud Computing, the Paas layer is experiencing a fast growth. According to the Market Monitor research report [98], PaaS accounted for the 24% of the total public cloud revenue in 2012 and it is expected to grow between 2012 and 2016 at a 41% compound annual growth rate (CAGR). Many companies are behind this success, in the following lines we describe some of them. Heroku [41] is a cloud Platform as a Service (PaaS). It has been in development since 2007, what makes Heroku one of the original PaaS offerings in the world. Nowadays it is owned by Salesforce.com since 2008. At the beginning, Heroku supported only the Ruby programming language, but since then, Heroku has been adding support for more programming languages: Java, Scala, Python, Node.js, Clojure, Grails, Gradle, Play and PHP. Heroku platform is entirely based on the AWS EC2 and S3, giving it the ability to scale in/out to satisfy customers’ demands. According to former Heroku CEO Byron Sebastian, Heroku was hosting more than 1.5 millions applications by November 2010 2. Openshift [69] is another interesting and new cloud PaaS product offered and developed by Red Hat. OpenShift is free and open-source, although now it has a paid version which adds extra support and features. The software that runs the service is called OpenShift Origin 3and can be downloaded from Github, allowing users to change it according to their needs.OpenShift is aimed at Java, Python, PHP, Node.js, Perl and Ruby developers and, as many others PaaS, OpenShift is based on AWS, this one specifically in Amazon EC2. CloudBees [25] offers a Java-based PaaS to host, run and manage Java applications. It is one of the first PaaS aimed mainly at the Java developer. CloudBees supports any JVM-based programming language or framework. Jenkins Continous Integration (CI) tool is included in the CloudBees PaaS. 2Heroku Boss: 1.5M apps, many not in Ruby - http://www.gigaom.com/2012/05/ 04/heroku-boss-1-5m-apps-many-not-in-ruby/ 3OpenShift Origin source code - https://github.com/openshift CHAPTER 1. INTRODUCTION 19 It supports developers through the whole application life cycle directly in the cloud from Github. The CloudBees service provides middleware on top of some public cloud, such as Amazon Web Services, OpenStack and VMware vSphere though customers can also run the service on private cloud infrastructure. 1.2.3 SaaS Software as a Service layer also strikes strongly in the Cloud race. SaaS represented 25% of total cloud revenue in 2012 and it is expected to follow growing in the following years [98]. In the next paragraphs, we highlight three successful SaaS companies without forgetting Google SaaS products Google Apps and Windows SaaS offerings previously described. Talking about successful SaaS solutions, we always find Dropbox [27]. Dropbox is a SaaS solution developed by Dropbox Inc., and launched on September 2008, that offers cloud storage and file synchronization to users. Dropbox uses Amazon S3 to store all files. In the released fact sheet of March 2012 4, Dropbox stated that they had over 50 million users, and nine months later, in November 2012, Dropbox announced that they had over 100 million users 5. Another notable SaaS is Salesforce.com [71], which is one of the most popular SaaS Customer Relationship Management (CRM) platform. It offers on-demand CRM services for all kind of organizations. Salesforce.com presents two major products. Sales Cloud is a set of applications to manage sales, customers and other business activities more easily and efficiently; and Service Cloud, which provides organizations with a community help-desk. Unlike Dropbox or many other SaaS solutions, Salesforce.com is built on its own infrastructure: Force.com [70], which is the Salesforce.com’s PaaS product. Talking about successful SaaS solutions, we have to talk about photo Cloud-storage. Nowadays, photo storage has become one of the most consumed SaaS solutions around the world. In this field, Flickr is one of the leading solutions. Flickr is an image and video hosting developed by Ludicorp in 2004 and adquired by Yahoo in 2005. According to a Verge article 6, in March 2013 Flickr reached 87 million members and more than 3.5 million 4Dropbox Fact Sheet - https://www.dropbox.com/static/docs/ DropboxFactSheet.pdf 5Dropbox Thanks a (hundred) million - https://blog.dropbox.com/2012/11/ thanks-a-hundred-million/ 6Verge report - http://www.theverge.com/2013/3/20/4121574/ flickr-chief-markus-spiering-talks-photos-and-marissa-mayer CHAPTER 1. INTRODUCTION 20 new images per day. It is written in PHP and use MySQL sharded cluster as it storage system. 1.3 Big Data Until now, we have showed the three most popular cloud paradigms: Infrastructure as a Service (IaaS), Platform as a Service (PaaS), and Software as a Service (SaaS) and how Google, Amazon, Microsoft and Rackspace among other companies offer their Cloud Services. Now it is time to move to one problem that has born with the outbreak of the Cloud Computing. Thanks to the features that Cloud Computing has brought with itself (pay-on-demand, elasticity, etc), lot of new applications have seen the light. Thanks to the new cloud business, they are now economically viable and before not. More and more novel applications harness the Cloud Computing paradigm every day, which means a never seen increase in the amount of generated as well as consumed data, called Big Data. Thus, scalable Database Management Systems (DBMS) have become a fundamental and critical part of cloud infrastructures [3]. Figure 1.2: Growth of data from the beginning of 2010 to 2020. Figure 1.2 depicts how fast generated data is growing 7. International Data Corporation (IDC) estimates the digital universe will grow up to 40,000 7Figure extracted from IDC’s Digital Universe Study, December 2012. CHAPTER 1. INTRODUCTION 21 exabytes, in more understandable words, 40 trillion gigabytes by 2020. From now until 2020, the generated data will double every two years [30]. Google has been one of the first companies that has addressed the Big Data problem, which is handling large amounts of data. Google has designed its own MapReduce programming paradigm [26], allowing them to do scalable distributed batch processing of large amounts of data: Web request logs, crawled documents, etc. Related to this field, Google designed and implemented their own distributed file system to be used in combination with their MapReduce framework, which was called Google File System (GFS) [32]. These are systems designed to run on a large cluster of commodity machines and are highly scalable, but both implementations remain private to Google 8. After the publication of MapReduce and the GFS, Apache Hadoop [77] entered the scene as the open-source alternative to Google MapReduce and Google File System. Apache Hadoop is an Apache project that includes implementations of a distributed file system, namely Hadoop Distributed File System (HDFS) [73, 94], and Hadoop MapReduce [77], both inspired by Google projects. Hadoop was initially created in 2004 by Doug Cutting (and named after his son’s toy elephant). In January 2008, Hadoop became a top-level Apache Software Foundation project. Nowadays it has many contributors, both academic and commercial (Yahoo being the largest commercial contributor), and continues growing. The field of distributed systems for Cloud Computing continued expanding, and as a logical next step, researchers needed a database for the applications to store the massive amount of data they were generating. Until then, the traditional and massive used Relational Databases Management Systems (RDBMS) have offered a simple and good solution according to the needs of the moment, but with the arrival of Web applications with massive numbers of users, the requirements of storage database systems for this new generation of applications have changed. Traditional databases did not offer a suitable and feasible storage solution to the new massive amount of data. In 2006, Google published a paper talking about BigTable [17]. Bigtable is a distributed storage system built on GFS for managing petabytes of structured data across thousands of commodity servers. Between its goals, we can find wide applicability, high performance, high scalability and high availability. To achieve such goals, BigTable uses a simple data model that supports dynamic control over data layout and format. Developers do not have to define a schema to store structured data, giving them a high flexibility when 8Colossus is the new version of the GFS as mentioned in the Spanner paper on OSDI 2012 [24]. CHAPTER 1. INTRODUCTION 22 building applications. Albeit, such a simple data store brings with it a lack of characteristics. Lacking in particular are ACID transactions, Join operations and accessing through a natural query language like SQL that RDBMSes offer. As a response to BigTable we find Apache HBase [86], which is an open source project, modeled after Google’s BigTable and written in Java. Developed as part of Apache Software Foundation’s Apache Hadoop project, it became an Apache top-level project on 2010. Hbase uses HDFS as the underlying storage system for the created tables and the Apache ZooKeeper as a distributed coordination service, similar to the use of Chubby [15] in BigTable. Hbase features are similar to BigTable features; its implementation is very close to BigTable implementation with same properties. Main differences lie in implementation details. For example, how the memory is mapped or how the Garbage Collector works [72]. HBase, BigTable and many others Cloud data stores are included in the group of so-called NoSQL databases 9. All of them are distributed storage systems mainly designed to offer a really good performance when dealing with Big Data, contrary to what traditional SQL databases do. This thesis is structured as follows: In the next Chapter we present background knowdledge of Cloud-based data stores and more in detail HBase data store. In Chapter 4 we characterize the environment we have been working in. We describe the cluster used for testing , the software deployed, and the dataset of our experiments. In Chapter 5 we discuss about the steps taken to import our dataset into a fully tuned HBase cluster and compare obtained results. In Chapter 6 we discuss about data retrieval in our HBase cluster and show the outcomes. In Chapter 7 we compare our HBase cluster against a MySQL cluster and the results are commented. Finally, in Chapter 8 we discuss about opportunities for improvement and review the work done for the Final Project. 9List of NoSQL databases - http://nosql-database.org/ Chapter 2 Background 2.1 Datastores: From SQL to NoSQL systems Adatabase is an organized collection of data items, which are records of some real world information [39]. Databases have always been extremely important in our society, from time ago with non-digital databases to nowadays with digital ones. They are a ubiquitous part of today’s computing environment. These database systems follow some data model, whose purpose is to determinate the logical structure of data items and how they are stored, organized and manipulated in data structures. As a vastly used data model we can found the relational model, used in SQL-based databases. This data with its data model needs an structure to rely on, and this is called Database Management System (DBMS). A DBMS is a suite of computer software programs that provides an interface between users and a database. They define mechanisms to build, store, maintain and modify one or more databases. DBMSes support any kind of applications, from business to Internet applications. They are one of the most important parts of many organizations and run critical applications that hospitals, airlines, banks and other types of organizations rely on for their daily operations. A relational database is a database which uses the relational model as its data model [21]. Its DBMS receives the name of Relational Database Management System (RDBMS) and it is the most popular example of database model. Most RDBMSes employ the SQL data definition and query language. Over the last three decades, RDMBSes have been the main technology for storing structured data. Even nowadays, the most popular DBMS continues being relational DBMS [74]. RDMBS have been proved to be a good solution and have been evolved to fit new application requirements. These relational 23 Chapter 3 Technical background In the following section we focus on the Column-oriented datastores, explained before, and in particular in the HBase solution, one of the most popular and open-source Column-oriented datastores. 3.1 Column-Oriented Datastores When referring to Column-Oriented datastores, also called Extensible Record Stores, where BigTable is the pioneer. BigTable and many other datastores of this type present a simple and flexible data model that can be extended at any moment. Some famous Column-Oriented datastores are Apache HBase [86], HyperTable [44], Apache Accumulo [79], and Apache Cassandra [78] in addition to many others. Among all these solutions, HBase [31] is likely the most popular opensource Column-Oriented datastore and is the one that we are going to use for our experiments. In the next section, we will describe HBase and how it works in order to fully understand this Thesis. 3.2 HBase HBase is an important Apache Hadoop-based project, which, as we stated before, is modeled on Google’s BigTable database. HBase can be characterized as a distributed, fault tolerant scalable database built on top of the HDFS file system. It belongs to the group of column-oriented datastores and uses Apache ZooKeeper for management of partial failures. Below, all must-known aspects of HBase are presented: HBase data model, storage, architecture, write/read/delete paths and the client API. 30 CHAPTER 3. TECHNICAL BACKGROUND 31 3.2.1 Data Model HBase stores data items, those are key/value pairs. The keys are multidimensional. Each single value is indexed by a row key, column key and a timestamp 1. Row keys are unique and allow the user to address all columns in one logical row. The column key is the combination of a column family and a column qualifier. Column families are the main unit of separation within a table. The last part is the timestamp which is used for versioning the data. Timestamp is usually automatically generated by the corresponding Region Server but it can be specified by the user. Summarizing, a key is commonly represented as the tuple (row, column, timestamp), which addresses a specified value: (row, column, timestamp) ->value where column equals (column family, column qualifier), timestamp is a 64-bit integer and row, column family, column qualifier and value are uninterpreted array strings, since in HBase everything is store as bytes. Row Key Time Stamp Family:Qualifier Family:Qualifier ”com.cnn” t9 anchor:cnnsi = ”Y” ”com.cnn” t8 anchor:look = ”X” ”com.cnn” t6 contents:html = ”” ”com.cnn” t5 contents:html = ”” ”com.cnn” t3 contents:html = ”” Table 3.1: Example HBase table given by [83] A cell is a set of data items with a common row and column key (remember that column key stands for (column family:column qualifier)), being the cell key = (row, column). Each data item in a cell is called a version of that cell. HBase supports multiple versions of cells. Each version of a cell is stored as a separated cell, next to other versions of that cell. Versions of a cell are sorted descending by the timestamp so that users will see the newest value first while reading it. Another HBase table feature is that it does not store NULL values as RDBMSs do. Files storing the data only contain data explicitly set. A table is organized by grouping cells into rows. These cells are sorted lexicographically by row key first and then by column key, so row keys lexicographically close will be stored near to each other. The sorting allows the 1Timestamp, also known as version number. CHAPTER 3. TECHNICAL BACKGROUND 32 table to be partitioned into the denominated regions, which hold exclusive ranges of row keys. The regions of a table are distributed between different nodes. 3.2.2 Storage The real view of tables differs from the conceptual view explained before. Physically, tables are stored on a per-column family basis, which means that all column family members are stored together in files called HFiles/StoreFiles. Such an approach brings advantages, one of which is that they can be compressed together. Column families must be declared while creating the table, whereas column qualifiers can be added to column families at any time. As written before, the HBase storage files are called HFiles. They are based on Hadoop’s TFile2class and mimic the SSTable format used in Google’s BigTable system. What is stored inside them is called KeyValue instances. They are the physical view of the conceptual cells. The whole cell, with its structured data (row length, key type, etc.) is what is a KeyValue object. Figure 3.1 shows a conceptual depiction of the KeyValue format. Figure 3.1: The KeyValue format, extracted from HBase: The definitive guide [31]. 3.2.3 Architecture HBase consists of three layers: the client, the server and the storage layers. The server layer consists of a master server and many region servers; the client has the library to communicate with the existing HBase installation; the storage layer is composed of a file system and a coordination service. Hadoop Distributed File System (HDFS) is the most used and tested file system to work with HBase. As the coordination service, ZooKeeper is the one HBase uses for its distributed coordination service. In the following section each component is explained as well as some HBase’s features. 2TFile Specification - https://issues.apache.org/jira/secure/attachment/ 12396286/TFile+Specification+20081217.pdf CHAPTER 3. TECHNICAL BACKGROUND 33 Figure 3.2: HBase architecture overview 3.2.3.1 Storage layer The storage layer is composed of a chosen file system for the HBase cluster and a coordination service: •File system: HDFS [12, 77] is the default file system when deploying an HBase cluster, but optionally it can be replaced by any other file system. For HBase, HDFS is the primary option as it is scalable, fail safe, has automatic replication and is built to run on commodity hardware. It fits with the needs of a distributed system. •Coordination service: Apache ZooKeeper [43, 48, 87] is an open source project and a part of the Apache Software Foundation. ZooKeeper is a highly-reliable distributed coordination service, comparable to Chubby [15], owned by Google and used for BigTable. Its aim is to offer a file-system-like access to clients with directories and files (znodes) that are used to store data, register services or watch for updates in a simple interface. In HBase, every region server creates a node in ZooKeeper, which will be used by the master to discover them. HBase uses ZooKeeper to choose the unique master and to store ”-ROOT- ”’s address too. Master server and region servers communicate with ZooKeeper to keep track of the current situation of the regions and region servers. CHAPTER 3. TECHNICAL BACKGROUND 34 3.2.3.2 Server layer Server layer is compound of two parts: •Master: The master server does not store any actual data and is not part of the retrieval path. It is responsible for assigning regions to region servers, handling load balancing of regions across region servers, unloading busy servers and moving regions to servers that are freer, and performing garbage collection of files. Since it never stores or provides data to clients or region servers, is usually slightly loaded. Furthermore, it takes care of all administrative operations such as schema changes or creation of tables or column families. If the master server goes down, the cluster can still work as the HBase client does not talk with it directly. Nevertheless, master should be restarted as soon as possible. •Region: As Lars George states in his book HBase: The definitive guide [31], a region is the basic unit of scalability and load balancing in HBase and is responsible for storing the actual data. Inside it, we can see contiguous ranges of rows stored together. Clients communicate directly with region servers to run all retrieve / write data operations. Each region is served by only one region server, although each region server can store multiple regions. The region servers also split regions into two pieces at the row key which is in the middle of the whole region once they have exceeded the maximum allowed region size. 3.2.3.3 Client layer Client needs to be able to find the region server which has a specific key. In order to achieve that, client communicates with the Zookeeper server and then retrieves the location of a table called ”-ROOT-” table from there. This table stores information about all regions in the ”.META.” table, which is another table storing where regions are and its row key ranges. So client reads the ”.META.” table and receives the exact address of the correct region handling a determined key or key range. Thanks to this three level lookup, the client is able to find the correct region server and perform its operations. For optimizing the process, client caches region locations because once the user table region is known, it can be accessed straightforwardly without the three level lookup. Only if the region server informs the client that it is no longer serving a region, does the client a new three level lookup. CHAPTER 3. TECHNICAL BACKGROUND 35 3.2.4 Write This subsection clarifies how writes and related stuff are done in a such a complex system as HBase is. Write requests are served by Region Servers. These components have three main elements: Write-Ahead Log (WAL), Memstore and HFiles. The WAL acts as a log for all modifications done to data in that exact region, guaranteeing atomicity and durability. The WAL allows regions to recover from server failures. It contains all previous modifications and can be used to recover from and to replay the data, achieving the last known stable stage. Memstore is an in-memory buffer that contains recently updated data items sorted by key. It works similar to how write buffers work on microprocessors, MemStore buffers writings, thus reducing write latencies. There is one MemStore for each column family of a table in the region. HFiles stores actual HBase’s data. HBase stores data to disk in a way similar to how a Log-Structured Merge Tree works [61]. First, when new data arrives at the region server it is put on the Write-Ahead Log (WAL). If the write to the WAL is achieved, the region server writes the data into the memory store. Region continues writing data to the Memstore until a configured maximum size. Once that threshold is reached, Memstore flushes the data to disk. When flushing, multiple HFiles are created, one per column family, which can affect negatively to HBase performance. Massive amounts of tiny files stand for more seeks and thus, higher latencies while reading. Due to that fact, HBase monitors the number and size of these files and does compactions from time to time. A compaction is the process of merging multiple HFiles together, and there are two types: •Minor compaction: It is responsible for rewriting the last few files into a larger one. It is triggered when certain properties and ratios are reached. •Major compaction: It merges all files of a region into one single file. It is usually triggered every twenty-four hours or even more due to heavy load. Sometimes, it is recommended to disable major compactions and manage them manually as they incur a high performance penalty due to rewriting all of the database contents. HFile files store an index at the end of the file to locate blocks within the HFile. The index is loaded into memory when the HFile is opened allowing look-ups to be performed with a single disk seek. For a more complete overview of how these files are designed refer to [94]. CHAPTER 3. TECHNICAL BACKGROUND 36 3.2.5 Read This paragraph explains reads in HBase. First of all, when a Region Server receives a Get request, it checks whether the desired row is in the MemStore or not. If not, it starts to search through the HFiles starting from the newest working towards older HFiles, which means to take a look through the disk contents because the data could be spread over multiple files. 3.2.6 Delete Here we explain how data is deleted from a HBase cluster. We must be clear that rows are never directly erased from HBase. When a Region Server receives a Delete request, it looks for the row and writes a delete marker to it. Whenever a Get request tries to access a row that has been previously deleted, it will find the delete marker and data will not be returned. During the next major compaction, rows with the delete marker will finally be deleted. 3.2.7 HBase API HBase provides a powerful client API written in Java as HBase does. It provides from basic operations to expensive ones. The initial set of basic operations are called CRUD operations which stands for ”Create, Read, Update and Delete”. Put Operation There are two groups of Put operations, first one works on single rows and the other works on a list of rows, both allow users to store data into HBase’s tables in a transparently manner. Defined as: put([List]Put put[s]) The user needs to supply a Put object or a list of them. These Put objects are created with the Put constructor: Put(byte[] row) A row key is supplied in order to create the Put instance, and once the user has created it, he/she can start to add values to the specified Put instance with the Put.add(value) method. The user can also supply a version number for a given key/value pair (timestamp), but if it is not specified, HBase gives to it a version number created from the current time of the Region Server responsible for that given row. CHAPTER 3. TECHNICAL BACKGROUND 37 When providing the value for the Put.add() method, as opposed to what happens with column qualifiers that can be whatever the user needs, an existing column family need to be given. Column families are usually defined when creating tables, but the user can always add new families calling to an expensive operation. It is because of that, HBase heavily recommends users to use fixed column families, although they can be altered. Unlike column families, new column qualifiers, version numbers and row keys can be provided on-the-fly within Put operations with no extra cost. The HBase API allows to use a built-in client-side write buffer that collects sets of Put operations before sending them as a unique RPC connection to the corresponding server. Hence, less RPC connections are needed resulting in an increase of the overall performance. HBase is smart enough to group and sort Puts by Region Server. If write buffer is not used, any time a Put operation is completed, HBase’s API will submit it to the right Region Server. Atomic Compare-and-Set Operation As a variation to the Put calls, the user can use Check and Put operation, defined as: checkAndPut(byte[] row, byte[] family, byte[] qualifier,byte[] value, Put put) This method allows users to issue Puts with a checking point. Only if the check step is successfully completed, the put operation is performed, everything as an atomic operation. It is really useful when dealing with data that needs previous values or similar stuff. All rows inside the Put object must be equal to the given row. The user can not use this operation to check different row keys. Otherwise, Check and Put operation will fail. Get Operation Like in HBase Put method, there are two groups of Get operations, the ones that work on a single row and the others that operate with multiple rows. Both allows the user to retrieve data stored in HBase’s tables. Get operation is defined as: get([List]Get get[s]) The user needs to supply a Get object. These objects are created with the Get constructor like in Put operations. it is: Get(byte[] row) CHAPTER 3. TECHNICAL BACKGROUND 38 In analogy with Put constructor, the user provides a row key to the Get method in order to get a Get instance. Get operation is bounded to a specified row, but can retrieve any data stored in it, from one value to all columns with its values. The user can add parameters to the Get object in order to narrow down the search. They will act as filters. If the user wants everything of a row, no filters are used, but if the user only wants a column family, the Get.addFamily() method must be used. Same thing happens if the user only wants a column qualifier from a column family (Get.addColumn()), or the row with a known timestamp (Get.setTimeStamp()). Lastly, there are methods, acting as filters, that allow users to specify how many versions want to be retrieved (Get.setMaxVersions()) and many others. Like in Put method, the user can retrieve a list of Gets, instead of only one row. The main difference is that the user issues a list of Gets, instead of one Get object. The result will be an array of Results, one for each Get instance. Delete Operation HBase client API provides a method to delete data from its tables. Delete method is defined as: delete([List]Delete delete[s]) Once more, the user is able to delete one row by one row or a list of them, the difference is the type of parameter: a Delete instance or a list of them. As with Get and Put calls, the user has to create a Delete instance and then adds details (filters) about the data he/she wants to remove, the constructor is: Delete(byte[] row) A row is provided. Subsequently, what user wants to be removed is added to it using different methods. Most important ones are Delete.deleteFamily() method, used to remove an entire column family, including all its columns, and Delete.deleteColumns() method, which operates on one column of a given column family, deleting all versions contained or just the cells matching the timestamp if it is provided. There are other types of Delete methods, but less used. Atomic Compare-and-delete Operation As a variation to the Delete call, the user can use Check and Delete operation. It is defined as: CHAPTER 3. TECHNICAL BACKGROUND 39 checkAndDelete(byte[] row, byte[] family, byte[] qualifier, byte[] value, Del del) It works as a Delete operation but adding a previous step in which a specify row key, column family, column qualifier and value are checked before deleting the desired row. Users can only check and delete on the same row. If the row key differs from the one pointed by the Delete instance, the CheckandDelete operation will fail. Useful hint. Row Locks for Row mutations: Previous operations: Put, Delete and CheckAndDelete are executed in such a way that they guarantee row level atomicity. They are executed entirely. A row lock is provided by the corresponding Region Server, protecting the row from other users trying to access it. Not from users trying to read it, only for those submitting row mutations operations. Scan Operation Scan operation allows the user to scan a range of data, from one row to a determinate stop row, taking advantage of the underlying sequential storage layout HBase has. The scan returns every row between the chosen range of rows. This operation is not executed atomically, it can be partially executed. Hence, returned data can be outdated if there has been a write operation during the scan. Scan operation is really similar to Get method and it works as a iterator, it means that user has a Scan instance and he/she must iterate over it to get all the results. Scan method is defined as: getScanner(Scan scan) getScanner(byte[] family) getScanner(byte[] family, byte[] qualifier) Each of them narrow the read data, from a general Scan to one that only returns values from a Column Family and inside it, from a Column Qualifier. As in Put or Get methods, Scan object is created with its constructor, which is: Scan(byte[] startRow, byte[] stopRow) It returns a Scan instance. Start row is mandatory and always inclusive, while stop row is not and is exclusive ( [startRow, stopRow) ). The user can submit filters as well. CHAPTER 4. ENVIRONMENT 46 sub-list of N sub-elements. While the total size of the XML is known, the amount of elements within each document is unknown due to the unfixed size of the XML files. The same situation happens with the size of each element where its number of sub-elements is unfixed as well as the length of each one (one can be a string sentence really long while other an integer). A conceptual example of one of our XML files is shown below 1. 1<element 1> 2 . . . 3</ element 1> 4<element 2> 5<sub−element 1> 6 ”Hi , I am a looooooong s t r i n g ” 7</sub−element 1> 8<sub−element 2> 9 ”Hi , I am a looooooong s t r i n g v ersion 2” 10 </sub−element 2> 11 . . . 12 <sub−element n> 13 <sublist1> 14 . . . 15 </sublist1> 16 </sub−element n> 17 </ element 2> 18 . . . 19 <element N> 20 . . . 21 </ element N> 4.5.1 HBase storage schema design In the following chapters we will work with HBase as our datastore. For this reason, we have designed an HBase storage schema that is able to map our XML data previously depicted to our HBase database. The sparse nature of HBase tables (not all columns populated in a row) makes them an interesting storage substitute for our XML dataset, in which elements can have different number of items, or even different items. It is worth to state that this schema is based on our main data access pattern in order to support efficient performance when updating and retrieving data from it; our row key is the Uid tag since most of the request we will have to cope with will be Uid-based requests. Nonetheless, there could be retrievals of some other nature, but they will only represent a low percent of the total requests. Besides the row key, there are four column families which 1For the sake of simplicity, only a basic form of a real XML is depicted. CHAPTER 4. ENVIRONMENT 47 map all the other sub-elements within our XML elements. ”Main” column family maps the main sub-elements of each XML element while the other three column families map sublists within each XML element (See conceptual XML example above in order to understand what are ”element” and ”sublist” words). Chapter 5 Design and Implementation - Import In the next two chapters, we will discuss about the steps taken to import our dataset into a fully tuned HBase cluster deployed on top of Triton. We will go through a basic version of importing data to the most fine-grain solution we managed to get, but before going into details, let us depict the work-flow taken, doing your reading more satisfactory and easier. Figure 5.1: Import research workflow. 48 CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 49 5.1 HBase cluster at a glance We use a 5 nodes HBase cluster by default: 1 Master server and 4 Region Servers. Default parameters are described below. Whether there is a change in the number of nodes or any parameters, we will state it in its corresponding section. As we exposed in Chapter 4, we use HBase in combination with HDFS as our distributed filesystem. Hence, we deploy a DataNode in the same node where HBase Master is and a DataNode along with each Region Server. Besides it, if the task requires Hadoop MapReduce, we turn it on, starting a JobTracker where the NameNode is and as many TaskTrackers as DataNodes there are in the cluster. 5.2 HBase: Tuning parameters for a writeheavy cluster In this section, we will explain how we have optimized our HBase cluster to meet our needs, which could be summarized into ”excel in importing data” operations. HBase is highly configurable when it comes to data-writing with plenty modifiable parameters. In the following lines we specify which ones, why, and how we have modified them. They are Java Virtual Machine parameters, MemStore parameters, and a few Hadoop parameters. 1. JVM related parameters: -HBASE HEAPSIZE: The maximum amount of heap to allocate expressed in MB. We have increased this parameter from 1000 to 2000 as HBase is a RAM consumer. The more RAM, the better the performance. -HBASE OPTS: We have enabled Java’s garbage collector logs as a way to help us to discover how to improve performance by tuning JVM flags and looking for long and short pauses. Following Todd Lipcon blog articles ”Avoiding Full GCs in Apache HBase with MemStore-Local Allocation Buffers: Part 1, 2 and 3” 1, 1http://blog.cloudera.com/blog/2011/02/avoiding-full-gcs-in-hbase-with-memstorelocal-allocation-buffers-part-1/ CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 50 we have turned on the Parallel New collector for the young generation (-XX:+UseParNewGC) and the Concurrent Mark-Sweep collector for the old generation (-XX:+UseConcMarkSweepGC). The Parallel New collector is a ”stop-the-world copying collector” but since the young generation is small and it uses many threads, the collector finishes its work very quickly and no stops are apparent. The Concurrent-Mark-Sweep collector (CMS) is responsible for cleaning dead objects in the old generation. It is also a ”stop-the-world collector”. The problem is that sometimes CMS fails and pauses of more than a minute appear in logs. CMS has two failure modes: (a) Concurrent Mode Failure: To avoid it, we need the garbage collector to start its work earlier in order to avoid getting overrun with new allocations. Setting -XX:CMSInitationOccupancyFraction flag to 70 turns out to help us. (b) Promotion Failure Mode due to fragmentation: This happens when there is not enough contiguous free space in the oldgeneration to allocate objects. This is termed memory fragmentation. When this occurs, the copying collector is called owing to its ability to compact all objects and free up space. To address this issue and avoid the stop produced by the copying collector, we use MemStore-Local Allocation Buffer (MSLAB) 2, a new Todd’s experimental facility. Hbase.hregion.memstore.mslab.enabled flag is set to true and the hbase.hregion.memstore.mslab.chunksize is set to 2MB per memstore. hbase.hregion.memstore.mslab.enabled = true hbase.hregion.memstore.mslab.max.allocation = 256KB hbase.hregion.memstore.mslab.chunksize = 2MB Table 5.1: HBase MSLAB parameters. To get a deeper information about this two modes or how garbage collector and HBase work together, read Todd Lipcon GC blog article [53] and HBase Documentation Chapter 13 Troubleshooting and Debugging Apache HBase [85]. 2MSLAB articles for a deeper background [84] [54]. CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 51 2. MemStore parameters: HBase write operations are applied in the hosting region’s MemStore at first, and then flushed to HDFS to save memory space when MemStore size reaches a threshold. This is what happens in a normal write scenario, but in a write-heavy HBase cluster we may observe an unstable write speed because updates are being blocked by Region servers. There are three blocking scenarios: •Size of all MemStores in a region server reaches a maximum and all the updates are blocked and flushes are forced. •Region’s MemStore size reaches a threshold defined by memstore.flush.size *memstore.block.multiplier. •A Store has more than hbase.hstore.blockingStoreFiles number of StoreFiles (one StoreFile per MemStore flushed). To avoid update blockings due to write-heavy workloads we have tuned MemStore size and related parameters, such as upper and lower limits before flushing and blocking times, following the configuration parameters for an HBase heavy-write load cluster proposed on chapter 9 ”Advanced configurations and Tuning” of the book ”HBase Administration Cookbook” [45], experiences from Sematext [9] and GBif company [58] and the most important, following our studies about our own HBase logs: (a) hbase.regionserver.global.memstore.upperLimit set to 40% (default one) (b) hbase.regionserver.global.memstore.lowerLimit set to 35% (default one) (c) hbase.hregion.memstore.block.multiplier set to 8 instead 2. (d) hbase.hregion.memstore.flush.size set to its default value, which is 128MB. (e) hbase.hstore.blockingStoreFiles set to 20 instead of 7. This tuning has met our needs and therefore, it has allowed us to reduce the chances of update blockings. CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 52 5.2.1 Hadoop baseline Following HBase tuning parameters strategy, we explain how we have optimized our Hadoop system. Hadoop comes with a non-aggressive set of parameters by default which are proved to work well enough. Nevertheless, we are looking for the best performance and we must tweak them in order to be more aggressive. For our baseline Hadoop configuration we have made a few changes, exposed below: •The default number of map/reduce slots is not adequate for our workload, that is why we have modified it to run a maximum of 12 (instead of 2) simultaneously map tasks and 6 (instead of 2) simultaneously reduce tasks per node since our nodes have 12 cores each one. Always a bit over the total amount of cores per node. •mapred.child.java.opts Hadoop parameter caps the heap of each map/reduce task process at 200 MB, which is too small. We have overridden it to 3072MB. Mapred.child.ulimit parameter has been also modified to be 2.5 times higher than the new heap of map/reduce tasks to prevent out of control memory consumption. •dfs.datanode.handler.count controls the number of threads serving data block requests from Datanodes. We have set it to 8 instead 3 by default. Increasing this value will lead to an increase in the memory utilization of the Datanode, but since we have enough RAM, we can enhance it. •dfs.datanode.max.xcievers controls the number of files that a DataNode can service concurrently and it is commonly recommended to increment it from the default of 256 to something higher 3 4. We have set it to 512. •io.file.buffer.size parameter determines how much data can be buffered while operating with sequence files. We have raise it from 4096 to 65536 following Cloudera recommendation 4. •JVM reuse policy: mapred.job.reuse.jvm.num.tasks is a configuration parameter found in mapred-site.xml which decides wheter map/reduce tasks reuse or not spawned JVMs. We have set its value to -1, which means that an unlimited number of tasks can reuse the same JVM. This policy is expected to benefit in scenarios where there are many short-length tasks and this is exactly our case. 3http://blog.cloudera.com/blog/2012/03/hbase-hadoop-xceivers/ 4http://blog.cloudera.com/blog/2009/03/configuration-parameters-what-can-youjust-ignore/ CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 53 5.3 The import experiment Once we have described how we have tuned our HBaseHadoop cluster, it is time to focus on the first experiment itself, which is importing our whole dataset into a Cloud-based database like HBase and test how it works. We will go through a basic version of importing data to the most fine-grain solution we managed to get. Suffice to say that all obtained results and conclusions are disclosed and analyzed along the way. 5.3.1 First approach: An HBase client As first approach, we have developed a Java application which uses HBase Client API to import the whole data set into our 5 nodes HBase cluster. Basically, it creates the table that will hold the whole data and starts parsing the XML video files one by one. As we showed before, these XML files contain lots of elements. For each file, our client creates a list of Puts objects mapping each object to a parsed element, and subsequently, it is sent to the HBase cluster through a call to the HTable.Put API method. Once the XML file has been parsed, the application repeats the same process with the next file until all of them are read. Results: Total Elements imported 12186983 Elements/Sec 1508 Table 5.2: First solution: Results. We can see some flaws to this idea. Let’s explain them in order to understand our next approaches to the solution: 1. admin.createTable(): This method creates a table with only one region. This is an issue as HBase Client API is only able to communicate and send its Puts to only one region/node. While one node is taking all the work load, the others are idle. This behavior changes once a threshold is reached and the region is split into two halves by the RegionSplitter, and the Hbase Load Balancer enters in the scene distributing new regions across CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 54 the nodes, but until it is triggered, no more nodes are in use and the obtained performance is really poor. To understand this feature we can reproduce Ted Yu’s explanation from his technical article ”Load Balancer in HBase 0.90” [97], where he explains how Load Balancer works: ”If at least one region server joined the cluster just before the current balancing action, both new and old regions from overloaded region servers would be moved onto underloaded region servers. Otherwise, I find the new regions and put them on different underloaded servers. Previously (in the older Load Balancer version) one underloaded server would be filled up before the next underloaded server is considered.”. We can wrap up that we will see an improved performance once the region gets split and the Load Balancer starts to work. Figure 5.2: HBase one active Region Server. 2. Region Server handler count: Region server keeps a number of running threads to answer incoming requests to user tables. To prevent region server running out of memory, this property is set to 10 by default which is a very low number unless you are using large write buffers with a really high number of concurrent clients. In our case, because our payload per request is low, increasing this number to handle more requests from the client would be beneficial as it will mean more accepted concurrent write requests. On the other CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 55 hand, setting it to a higher number will consume more Region server’s memory, but our cluster has enough to handle this peak. So in order to leverage it, we need not just to change hbase.regionserver.handler.count parameter which affects the server-side, but also to make the client to work concurrently by using threads. hbase.regionserver.handler.count = 10 3. Compressed data: In Hbase, it is well-known that using some form of compression for storing data may lead to an increase in IO performance, and thus in an increase in the overall performance [66] [18] [4] [82]. But at this point, we are not using any sort of data compression yet. We should exploit it to reduce the number of bytes written/to read from HDFS, to save disk usage and to improve the efficiency of network bandwidth. On the other hand, if we enable it, we will need to un/compress data so we will need some extra CPU cycles. It is simply trading IO load for CPU load. 4. setAutoFlush: HBase client API provides a built-in Write Buffer which allows to cache a group of Put/Delete objects on the client side, and flushes these objects to the Region Servers in a batch so that they are sent within one RPC call to the servers, instead of sending Puts one at a time like by default. Using it, all requested changes will wind up in the same Write Buffer and will not be sent until the Write Buffer is filled. The chief advantage of using it is the reduction in the amount of necessary RPC connections to transfer data from the client to the sever and back. In our application, which needs to store thousands of values per second, less RPC calls will mean less round-trip times (RTTs) to happen. Figure 5.3 provides the architecture of the Client Write Buffer5. To take advantage of this feature, we must setAutoFlush to false, instead of true by default. 5This figure has been obtained from HBase: The Definitive Guide CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 62 5.3.3.3 Third approach third version: Pre-creating regions Here is where we address one of our older issues already discovered in the first approach: The method admin.createTable() creates a table with only one region. Now, that we are using Bulk Load MapReduce feature, we can see how this issue continues here just by glancing at Hadoop logs. The HFileOutputFormat.configureIncrementalLoad method looks up the current regions for our table and finds one, that is why it configures only one reduce partition (one reduce partition per region). Only one reducer task will be spawned, while the rest of the nodes will stay idle. Figure 5.7: HBase one single Region Server. Looking at the figure 5.7, we can see how data ends up within a single region in one Region Server. If we create an HBase table with only one region, all clients will only be able to write out to the same region until it gets split and distributed across our cluster. The solution is to pre-create a table with the desired number of empty regions; Admin.createTable(table, startKey, endKey, numberOfRegions) method allows us to do exactly what we want. It creates a table with numberOfRegions regions and as first split the passed startKey and as last split the endKey. We have configured it to pre-create 24 Regions in order to match the number of total reduce slots our cluster allows and thereby to complete the job spawning only one single wave of reducers. Figure 5.8 reveals a slight improvement performance in the needed time to import the dataset into HBase. Despite the performance enhancement, a closer look at the TaskTracker logs reveals some issues. Albeit all nodes are working now, some nodes are CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 63 Figure 5.8: Execution time to import the dataset with different number of pre-created regions. working harder than others. The graph in Figure 5.9 depicts the region distribution across the four region servers. Node cn212 stores the 72.22% of the total data, making an uneven distribution of it across the Region Servers. Next graph in Figure 5.10, shows the sizes of the 24 regions created. As before, data has been uneven distributed: Region 1 stores the 70% of the dataset. TaskTracker logs uncover what is the problem. Some reducers are working with more than 10 times the amount of records others are dealing with, which translates to different reducer execution times. While some finish within less than a dozen of seconds, others takes more than 5 minutes. This happens because data’s keyspace is not evenly distributed. Admin.createTable(table, startKey, endKey, numberOfRegions) flaw is that it uses Bytes.split as the split strategy and it does not work efficiently with unevenly distributed data. All the regions are accessible in the keyspace, but since our keyspace is not evenly distributed, some reducers/regions does not receive almost any data, while others collects nearly all data. CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 64 Figure 5.9: Uneven region server distribution. 5.3.4 Fourth approach: Coping with skewed data What we have experienced in the last approach is called Skew in a MapReduce environment. Skew refers to a significant load imbalance and its causes have been widely studied [6] [26] [90]. Skew can appear due to computational load imbalance, characteristics of the user-defined operations or of the specific dataset or by hardware malfunction among others reasons. Skew from either cause is undesirable because it leads to longer job execution times and throttles cluster throughput. The original MapReduce paper [26] tackles this problem using speculative execution. Albeit this works well, it is not the best solution since it means repeating work already done. Balazinska et al. identified a specific type of skew, referred to as Data Skew [51]: It affects both keys and values in either mappers or reducers. They state that data skew occurs more often for the reducers because mappers mostly take the same-size blocks of input data. There are two sub-types of this skew: one caused by uneven data allocation; the number of key values for one task is much larger than the number of keys in the other partitions to cause an imbalance. A second one caused by uneven processing times; one task processes larger number of values than the other tasks. According to our Hadoop cluster logs, data skew happens in the reducer phase because almost all mappers take the same time to complete their tasks but not the reducers. A deeper insight into logs reveals some reducers taking significantly larger number of keys than the other reducers. This is what is causing the imbalance situation and is referred to as Reduce phase: Partitioning skew [50]. CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 65 Figure 5.10: 24 uneven regions. In our MapReduce job, map tasks outputs are distributed among reduce tasks via TotalOrderPartitioner, which partitions the map output into ranges of the keyspace, which correspond to the region boundaries of our HBase table created by the Bytes.split method. This is not adequate for our data because it is not evenly distributed. There are lot of duplicated keys and a big part of them are really similar, ending up in the same region. To cope with this problem, we have to somehow find a good partitioning function that ensures total ordering, like TotalOrderPartitioner does, and splits the data into equal partitions as well. Hadoop provides a partitioning function called InputSampler, which sample the input at random or what user choose to estimate what is the best way to partition. But since it samples the map input data, it does not fit our needs. What we need to sample is the map output, which will be the keys of our table. That is why we have developed a lightweight MapReduce Java-based tool which samples map output keys and gives us a file describing the best partition for our dataset. Subsequently, this file can be used in combination with TotalOrderPartitioner to know which key/value pairs to send to which reducers, or it can be used in combination with admin.createTable(table, splitPoints) method to create a table with the best split points for regions in MapReduce-HBase environments. This file will be able to evenly span the key space creating an even distribution of records across the reducers and to create regions with almost the same size, therefore having a well apportioned HBase cluster. Our sampling tool uses a wrapper input format that makes a record reader which passes few key/value pairs to the mapper. The rate at which key/value pairs are passed to the mapper can be modified according to user needs. In CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 66 order to obtain a significant sampling of the entire data, adjust it to ten has been tested to be valid enough for us. Ten gives a good speed/significantsample ratio. The mappers of the sampling job emit only keys, while the values are always null. In order to reduce the total amount of generated data, the XML files are not completely parsed, it just needs the ids of each element. Finally, our tool also overwrites the input format with a sampling reducer that emits the exact number of samples needed for the creation of the regions of the HBase table. Next table shows how fast the sampling tool gets his job done: Sampling input dataset 83 sec Table 5.5: Sampling time result. We have modified the old MapReduce job to accept the text file created by the sampling MapReduce tool and to create the HTable with the new and correct splits points. The maximum number of reduce tasks that will be run simultaneously by a task tracker is set to 6 (mapred.tasktracker.reduce.tasks.maximum). Hence we create 24 regions in our table. 4 regions per node (6 simultaneous reducer tasks * 4 nodes = 24). Therefore, the job will only need one single wave of reducers to complete it. On the other hand, each map tasks will read off one DFS block, so multiple map waves will be used getting hide shuffle latency. Figure 5.11 shows the outcome of our tests. Using our Sampling Tool we have reduced the total execution time to import data to only 552 seconds, which is 34.29% faster than our previous results with pre-created Regions. Figure 5.12 depicts the region distribution obtained using the Sampling Tool. Now data is much evenly distributed along the four Region Servers. Digging into logs reveals more uniform reducers’ execution time as well. 5.4 Performance Tuning Hadoop At this point, we have reached the best possible importing performance level in our HBase cluster without going to deep into Hadoop parameters, so now we can start to fine tuning these configuration details in order to maximize the performance of our Hadoop workload. This tuning has been performed by taking the last approach as our baseline and following well-known studies about Hadoop performance tuning [8] [40]. In the following lines, we explain which parameters we have hacked and why: CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 67 Figure 5.11: Total time to import data with-/out using our Sampling Tool. •HDFS block size: In our Hadoop cluster, each mapper receives an input split whose size is determined by dfs.blocksize (by default, 64MB). If we increase it, the number of spawned mappers will decrease and less overhead will be created as there will be less map output splits to merge and less map tasks to run. On the other hand, the execution time taken by each mapper will increase. Figure 5.13 shows the performance with different HDFS block sizes. Our optimal size comes out to be 128MB. •Spilled records: While mappers are running, the generated intermediate output of map tasks is hold in buffers. Mappers have assigned a portion of memory of the map JVM heap in which they store their results, but if it gets completely filled up, its contents are spilled to disk. If this situation happens multiple times, it leads to additional overhead, which means more time to complete the phase. If we study our logs, we can see that the total Map output records is much lower than the Spilled records, which indicates that we are not setting an appropriate size for the buffers, they are being spilled to disk many times. To avoid this, we hack the value of the parameter io.sort.mb, which is by default 100MB to be big enough to hold all the records. By doing some calculations, setting it to 1280 MB fits our CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 68 Figure 5.12: Region distribution using our sampling tool. needs. Less records are spilled to disk and only the compulsory and final spill is done once the mapper is completed. If the input size of each mapper is 64MB: Our records have in average 3202.81 bytes/record so if the block size is 64MB we have 20953.12 records/input Every spilled record takes 16 bytes of metadata in buffer 20953.12 * 16 = 0.32MB 64MB data + 0.32MB metadata = 65MB needed. If the input size of each mapper is 128MB: Our records have 3202.81 bytes/record in average so if the block size is 128MB we have 41906.24 records/input Every spilled record takes 16 bytes of metadata in buffer 41906.24 * 16 = 0.64MB 128MB data + 0.64MB metadata = 128.64MB needed. We are still far away from the spill threshold setting. Same happens with the reducers, before applying the reduce function they need to copy, merge and sort the map outputs, so they start copying records from mappers and storing them in a buffer until a threshold is reached and then, these records are spilled to disk. The size of this buffer is governed by mapred.job.shuffle.input.buffer.percent parameter CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 69 Figure 5.13: HDFS block size setting. and its default value is 66% of the Reduce JVM heap space. The ideal scenario would be one where this buffer would be big enough to hold all map output records, but since it is a too high size or sometimes is even impossible to reach, increasing this percent to a higher number will be enough for our purposes. Finally, after running several tests, we saw performance improvements by increasing this parameter to 90%. Another related parameter is mapred.job.reduce.input.buffer.percent, set by default to 0%. It imposes the size of the Reduce JVM heap that is allocated to the final reduce function. Since our reduce function is not memory-bound, we can use a JVM heap percent to retain some records and thus reduce the number of IO operations. Consequently, we set it to 80%. The left side of the Figure 5.14 reveals that all the reduce input records (1780322 records) were spilled to disk within a given reducer task, while the right side of the figure shows that only 1238964 of a total of 1780322 records were spilled to disk. For the left side, the default settings described above where used, while for the right side, we hacked them whit the ones stated before. CHAPTER 5. DESIGN AND IMPLEMENTATION - IMPORT 70 Figure 5.14: Spilled records in reducer side. Chapter 6 Design and Implementation - Retrieval 6.1 Random reads in HBase Unlike some Cloud-based databases which are optimized for random reads like PNUTS, HBase is write-optimized by using on-disk structure that can be maintained using sequential IO. Its records are never overwritten, instead, updates are written sequentially to new files in disk. That means that multiple updates of the same record will be spread over many files, so when reading it, multiple IO operations will be needed to merge the separate updates. On the other hand, as we already explained, all writing is sequential, so HBase excels at writing and consequently, in scans, which are sequential reads. This is a simple trade-off between optimizing for reads and optimizing for writes. 6.1.1 Random reads in our heavy-write cluster In this section we test random reads for our HBase fully write-optimized cluster. In HBase, there is no big room for improving random reads, but still some improvements can be done to achieve a better random read performance than the default one. Before starting, we must describe how our reads are going to be, whether they will request an entire row, that will be the darkest room for enhancing it, or they will ask for a little part, better scenario as HBase stores familyColumns in separated files and only a few will be required in order to return the result. The use case we performance is fetching 1, 10, 25, 50 or more random video details at once. The row keys of these random elements are known beforehand, so we only have to look for them and retrieve its details. In 71 CHAPTER 7. DESIGN AND IMPLEMENTATION - BENCHMARKING78 simple sharding function S(key), S = hash(key) % numberOfNodes, because by default MySQL has no built-in clustering capabilities as HBase has. The MySQL table looks exactly like the HBase table does; ten fields by row, size of each one is the average size of our real data. InnoDB is used as the storage engine of our MySQL database and the row key is indexed by a B-Tree. DISCLAIMER: No MySQL-related parameters have been tuned for the benchmarks, but innodb buffer pool size increased from 8MB to 3027MB. It is a storage area for caching data and indexes in memory. Also a B-Tree index is created on the row key of our table to enhance the query execution time. The rest of parameters continue as by default. In the conducted benchmarks all fields are always read, the number of operations is always one million and the records to operate on follow a uniform distribution in order to be as close as possible to our real scenario fully described in the previous chapters. Table 7.1 summarizes the kinds of workload that we chose for benchmarking. Workload Insert % Read % Update % Data Load 100 Predominant Reads 95 5 Reads mixed with Updates 50 50 Table 7.1: YCSB Workloads Below we present results for each workload: load, predominant reads and reads with updates. 7.1.1 Load phase Along with the MySQL outcome, three different versions of HBase loads (see Figure 7.1) are depicted due to the lack of pre-split regions YCSB comes with. The first one, HBase label, shows the default behaviour of YCSB, which is a table with only one region at the beginning. The other two correspond to different split region algorithms we have tested. The first, HBase built-in PreSplit label, is based on the HBase.util.RegionSplitter tool which allows users to create tables with a specified number of pre-split regions, assuming keys are uniformly distributed bytes. This tool gives out a much better performance than the HBase version without pre-split regions. However, this can be improved a bit more because of the fact that YCSB does not export CHAPTER 7. DESIGN AND IMPLEMENTATION - BENCHMARKING79 Figure 7.1: Latency vs throughput comparison. uniformly distributed bytes keys. Instead, it creates keys which are a combination of the string user and a long integer (Ex. ”user111111111111111”). Once we understood the YCSB key-creation pattern, we developed a custom lightweight RegionSplitter tool which leverages the YCSB key specification and creates a table with pre-split regions whose split points fits with the keys that YCSB will randomly create during the load phase. This tool, HBase PreSplit Tool label, overcomes the previous results by achieving 20356 operations per second, 66.25% better compared to the YCSB default throughput behavior (HBase label). We conclude that our tuned HBase cluster has unconquerable superiority in writes while MySQL is really far away from the HBase results. 7.1.2 Reads mixed with Updates Figures 7.2 and 7.3 measures the latency/throughput curve of both HBase and MySQL clusters when dealing with a workload composed of 1.000.000 CHAPTER 7. DESIGN AND IMPLEMENTATION - BENCHMARKING80 Figure 7.2: Workload 50% read 50% update - Update. operations, 50% are update operations and the other 50% are read operations. HBase is optimized for writes and that is why it achieves higher throughput and lower latency than its competitor. The main reason of this performance is because edits are commited to memory firstly (WAL >>MemStore) and then aggregated edits are flushed to disk. 7.1.3 Predominant Reads Figures 7.4 and 7.5 measures the resulting latency/throughput curve of both HBase and MySQL clusters when dealing with a workload composed of 1.000.000 operations, 95% of them are random reads and the 5% rest are update operations. Upon studying the results we conclude that MySQL Sharded is a performance leader in reads. MySQL B-Tree indexes make the difference. Its only flaw is that it gets satured when the offered throughput reaches more than 9734 operations per second. Random read performance is slower in HBase because of the need to reconstruct records, but not much. Talking about Update, HBase results are from other world. MySQL has nothing to do against HBase when it comes to writes. CHAPTER 7. DESIGN AND IMPLEMENTATION - BENCHMARKING81 Figure 7.3: Workload 50% read 50% update - Read. Figure 7.4: Workload 95% read 5% - Read. CHAPTER 7. DESIGN AND IMPLEMENTATION - BENCHMARKING82 Figure 7.5: Workload 95% read 5% - Update. Chapter 8 Conclusions In this chapter we discuss opportunities for improvement in Section 8.1 and review the work done for the final project in Section 8.2. 8.1 Future work Dealing with data, no matter whether it is import or retrieval operation, has been carefully studied and discussed. Nonetheless, some improvements not tested arise here: •Followed scheme design for our HBase table has proven to work well. However, we may consider a redesign of it. A schema with only one columnFamily would be beneficial as we would have a better control over the HBase behavior; easier manage of storeFiles, reads/block caches and similar opportunities derivated from having only one columnFamily. •Dataset has lot of duplicates elements. By now, we just import them, it does not matter whether they are already stored or not. Nonetheless, we could create a combiner class to get rid of duplicates in the mapper side. It would reduce the amount of IO operations between mappers and reducers and would reduce the total execution time of the job. The drawback would be that only one version of each element would be stored in the database. •In HBase, it is possible for a client to read directly from disk instead of going through the DataNode. This action is called a short-circuit read. Region servers read directly off the local node data disks instead of asking the DataNode for the data. This feature has been tested to 83 CHAPTER 8. CONCLUSIONS 84 work well, with little or no drawbacks, hence we could use it instead the default HBase built-in read behavior. •Relationship between data disks and Hadoop / HBase ecosystems has been and continuous being the focus of a lot of research activity [49] [28] [7]. Researcher Shrinivas B. Joshi points out the advantages of using more than one data disk in Hadoop workloads (achieved more than 50 % performance improvement) [47]. It is well-known Hadoop performance scales with the number of available data disks, however, we were not able to check it owing to our hardware boundaries, but it may be worth trying it out. •It is well-known that HBase communication stack does not work correctly when using high performance networks like InfiniBand because of its implementation based on Java Sockets Interface that provides non-optimal performance due to the created overhead [92]. Although we are using Gigabit Ethernet for our experiments, Triton cluster provides InfiniBand network communication. Therefore, we could harness it by using the novel desing of HBase that Jian et al. have done for their research, a fully InfiniBand compatible HBase [42]. They claim to have achieved a factor of 3.5 improvement over 10Gigabit Ethernet network latency when retrieving data (Get operations). •HOYA1: It is a YARN application2that provisions Region Servers based on an HBase cluster configuration, it can be used to spin up temporary HBase clusters during MapReduce or other jobs. HOYA could be useful for our interests due to it would help us to spin up more HBase resources during heavy batch workloads such as night import of new data. It will allow us to create on-demand HBase resources and thanks to it we will be able to utilize cluster resources better. A framework for job scheduling and cluster resource management 8.2 Discussion In this final project we present methods related with scaling-out the data of a commercial company. In order to improve the performance we implement 1HOYA https://github.com/hortonworks/hoya/ 2YARN, called the MapReduce 2.0, is a framework for job scheduling and cluster resource management - http://hadoop.apache.org/docs/current/hadoop-yarn/ hadoop-yarn-site/YARN.html CHAPTER 8. CONCLUSIONS 85 an HBase cluster, a Cloud-based datastore, along with solutions based on Big Data algorithms, such as MapReduce. This project evaluates the obtained performance from three different points of view: Firstly, the main problem of importing a really big dataset to a new Cloud-based datastore. Several approaches have been developed and carefully tested, uncovering its benefits and drawbacks in order to improve the obtained approach. Secondly, the performance of reading random data in a write-optimized database like HBase. Once more, conceptual ideas have been developed and tested and the results have been exposed. Finally, the tuned HBase cluster has been benchmarked against a MySQL cluster similar to the one the company where the data comes from uses. We have been able to improve the default performance in every area. Importing data we have passed from an API client whose execution time is more than five hours to a MapReduce-based client which enables us to reduce the processing time to only nine minutes. There are lot of improvements behind these simple numbers, such as HDFS issues, skew data, compression, JVM issues, etc. As of retrieving random data, we have also improved default results by using concepts such as Bloom filters, HFile’s block sizes or block caches. All obtained results have been studied and improved when possible. Setting aside the results, one main tool has been developed not only to fit the uneven data issue our dataset has, but also every MapReduce/HBase job suffering from skew data in the mapper outputs. In a brief way, it samples the whole dataset in a lightweightly way with a confident level defined by the user and returns the best split points which removes the uneven distribution of mappers output. The conducted benchmarks shows how our tuned HBase cluster performs against a MySQL cluster. Three main scenarios are developed and the outcomes are discussed. HBase outperforms in writes and is really close to MySQL in random reads. It is worth to state we have improved the Yahoo! Cloud Benchmark Tool by developing some tools to overcome pitfalls already presented in its solution and thus letting us enhance HBase results (not hack them, but get closer results to real scenarios). We can conclude that we have achieved enough good results as to change the company datastore backend system to HBase. Bibliography [1] Abadi, D. Consistency tradeoffs in modern distributed database system design: CAP is only part of the story. Computer 45, 2 (2012), 37–42. [2] Agaoglu, E. LZO vs Snappy vs LZF vs ZLIB, A comparison of compression algorithms for fat cells in hbase, 2013. http://blog.erdemagaoglu.com/post/4605524309/ lzo-vs-snappy-vs-lzf-vs-zlib-a-comparison-of. Accessed 11.8.2013. [3] Agrawal, D., Das, S., and El Abbadi, A. Big data and cloud computing: New wine or just new bottles? Proceedings of the VLDB Endowment 3, 1-2 (2010), 1647–1648. [4] Aiyer, A. S., Bautin, M., Chen, G. J., Damania, P., Khemani, P., Muthukkaruppan, K., Ranganathan, K., Spiegelberg, N., Tang, L., and Vaidya, M. Storage Infrastructure Behind Facebook Messages: Using HBase at Scale. IEEE Data Eng. Bull. 35, 2 (2012), 4–13. [5] Amazon.com Inc. Amazon Web Services web page, 2013. http:// aws.amazon.com/. Accessed 11.6.2013. [6] Ananthanarayanan, G., Kandula, S., Greenberg, A. G., Stoica, I., Lu, Y., Saha, B., and Harris, E. Reining in the Outliers in Map-Reduce Clusters using Mantri. In Symposium on Operating Systems Design and Implementation (OSDI) (2010), vol. 10, p. 24. [7] Awasthi, A., Nandini, A., Bhattacharya, A., and Sehgal, P. Hybrid HBase: Leveraging Flash SSDs to improve cost per throughput of HBase. [8] Babu, S. Towards automatic optimization of mapreduce programs. In Proceedings of the 1st ACM symposium on Cloud computing (2010), ACM, pp. 137–142. 86 BIBLIOGRAPHY 87 [9] Baranau, A. Configuring HBase Memstore: What You Should Know, July 2012. http://blog.sematext.com/2012/07/16/ hbase-memstore-what-you-should-know. Accessed 11.7.2013. [10] Bernstein, P. A., and Goodman, N. Multiversion concurrency control theory and algorithms. ACM Transactions on Database Systems (TODS) 8, 4 (1983), 465–483. [11] Bloom, B. H. Space/time trade-offs in hash coding with allowable errors. Communications of the ACM 13, 7 (1970), 422–426. [12] Borthakur, D. HDFS Architecture, June 2012. http://hadoop. apache.org/common/docs/r0.20.0/hdfs_design.html. [13] Brewer, E. CAP twelve years later: How the ”rules” have changed. Computer (2012), 23–29. [14] Brewer, E. A. Towards robust distributed systems. In Proceedings of the nineteenth annual ACM symposium on Principles of distributed computing (2000), ACM, p. 7. [15] Burrows, M. The Chubby lock service for loosely-coupled distributed systems. In Proceedings of the 7th symposium on Operating systems design and implementation (2006), USENIX Association, pp. 335–350. [16] Cattell, R. Scalable SQL and NoSQL data stores. ACM SIGMOD Record 39, 4 (2011), 12–27. [17] Chang, F., Dean, J., Ghemawat, S., Hsieh, W. C., Wallach, D. A., Burrows, M., Chandra, T., Fikes, A., and Gruber, R. E. Bigtable: A distributed storage system for structured data. ACM Transactions on Computer Systems (TOCS) 26, 2 (2008), 4. [18] Cheng, P., and An, J. The Key as Dictionary Compression Method of Inverted Index Table under the HBase database. Journal of Software 8, 5 (2013), 1086–1093. [19] Cloudera, Inc. CDH web page, 2013. http://www.cloudera.com/ content/cloudera/en/products/cdh.html. Accessed 1.4.2013. [20] Cloudera, Inc. Cloudera web page, 2013. http://www.cloudera.com/. Accessed 1.4.2013. BIBLIOGRAPHY 94 [93] White, T. The Small Files Problem, February 2009. http://blog.cloudera.com/blog/2009/02/the-small-files-problem/. Accessed 22.6.2013. [94] White, T. Hadoop: The definitive guide. O’Reilly Media, Inc., 2012. [95] Wikipedia. Cloud Computing — Wikipedia, The Free Encyclopedia, 2013. Online; accessed 10-June-2013. [96] Wikipedia. Consistency model — Wikipedia, The Free Encyclopedia, 2013. Online; accessed 10-June-2013. [97] Yu, T. Load Balancer in HBase 0.90, April 2011. http://zhihongyu. blogspot.fi/2011/04/load-balancer-in-hbase-090.html. Accessed 15.7.2013. [98] Yulitza Peraza, G. Z. 451 Market Monitor Cloud Computing: Overview Report 2013. 451 Research 5 (2013).