Showing posts with label scalability. Show all posts
Showing posts with label scalability. Show all posts

Thursday, 27 February 2014

The WhatsApp Architecture Facebook Bought For $19 Billion

Source: http://highscalability.com/blog/2014/2/26/the-whatsapp-architecture-facebook-bought-for-19-billion.html

Rick Reed in an upcoming talk in March titled That's 'Billion' with a 'B': Scaling to the next level at WhatsApp reveals some eye popping WhatsApp stats:
What has hundreds of nodes, thousands of cores, hundreds of terabytes of RAM, and hopes to serve the billions of smartphones that will soon be a reality around the globe? The Erlang/FreeBSD-based server infrastructure at WhatsApp. We've faced many challenges in meeting the ever-growing demand for our messaging services, but as we continue to push the envelope on size (>8000 cores) and speed (>70M Erlang messages per second) of our serving system.
But since we don't have that talk yet, let's take a look at a talk Rick Reed gave two years ago on WhatsApp: Scaling to Millions of Simultaneous Connections.
Having built a high performance messaging bus in C++ while at Yahoo, Rick Reed is not new to the world of high scalability architectures. The founders are also ex-Yahoo guys with not a little experience scaling systems. So WhatsApp comes by their scaling prowess honestly. And since they have a Big Hairy Audacious of Goal of being on every smartphone in the world, which could be as many as 5 billion phones in a few years, they'll need to make the most of that experience.
Before we get to the facts, let's digress for a moment on this absolutely fascinating conundrum: How can WhatsApp possibly be worth $19 billion to Facebook?
As a programmer if you ask me if WhatsApp is worth that much I'll answer expletive no! It's just sending stuff over a network. Get real. But I'm also the guy that thought we don't need blogging platforms because how hard is it to remote login to your own server, edit the index.html file with vi, then write your post in HTML? It has taken quite a while for me to realize it's not the code stupid, it's getting all those users to love and use your product that is the hard part. You can't buy love
What is it that makes WhatsApp so valuable? The technology? Ignore all those people who say they could write WhatsApp in a week with PHP. That's simply not true. It is as we'll see pretty cool technology. But certainly Facebook has sufficient chops to build WhatsApp if they wished.
Let's look at features. We know WhatsApp is a no gimmicks (no ads, no gimmicks, no games) product with loyal users from across the world. It offers free texting in a cruel world where SMS charges can be abusive. As a sheltered American it has surprised me the most to see how many real people use WhatsApp to really stay in touch with family and friends. So when you get on WhatsApp it's likely people you know are already on it, since everyone has a phone, which mitigates the empty social network problem. It is aggressively cross platform so everyone you know can use it and it will just work. It "just works" is a phrase often used. It is full featured (shared locations, video, audio, pictures, push-to-talk, voice-messages and photos, read receipt, group-chats, send messages via WiFi, and all can be done regardless of whether the recipient is online or not). It handles the display of native languages well. And using your cell number as identity and your contacts list as a social graph is diabolically simple. There's no email verification, username and password, and no credit card number required. So it just works.
All impressive, but that can't be worth $19 billion. Other products can compete on features.
Google wanted it is a possible reason. It's a threat. It's for the .99 cents a user. Facebook is just desperate. It's for your phone book. It's for the meta-data (even though WhatsApp keeps none).
It's for the 450 million active users, with a user based growing at one million users a day, with a potential for a billion users. Facebook needs WhatApp for its next billion users. Certainly that must be part if it. And a cost of about $40 a user doesn't seem unreasonable, especially with the bulk paid out in stock.  Facebook acquired Instagram for about $30 per user. A Twitter user is worth $110.
Benedict Evans makes a great case that Mobile is a 1+ trillion dollar business, WhatsApp is disrupting the SMS part of this industry, which globally has over $100 billion in revenue, by sending 18 billion SMS messages a day when the global SMS system only sends 20 billion SMS messages a day.  With a fundamental change in the transition from PCs to nearly universal smartphone adoption, the size of the opportunity is a much larger addressable market than where Facebook normally plays.
But Facebook has promised no ads and no interference, so where's the win?
There's the interesting development of business use over mobile. WhatsApp is used to create group conversations for project teams and venture capitalists carry out deal flow conversations over WhatsApp.
Instagram is used in Kuwait to sell sheep.
WeChat, a WhatsApp competitor, launched a taxi-cab hailing service in January. In the first month 21 million cabs were hailed.
With the future of e-commerce looking like it will be funneled through mobile messaging apps, it must be an e-commerce play?
It's not just businesses using WhatsApp for applications that were once on the desktop or on the web. Police officers in Spain use WhatsApp to catch criminals. People in Italy use it to organize basketball games.
Commerce and other applications are jumping on to mobile for obvious reasons. Everyone has mobile and these messaging applications are powerful, free, and cheap to use. No longer do you need a desktop or a web application to get things done. A lot of functionality can be overlayed on a messaging app.
So messaging is a threat to Google and Facebook. The desktop is dead. The web is dying. Messaging + mobile is an entire ecosystem that sidesteps their channel.
Facebook needs to get into this market or become irrelevant?
With the move to mobile we are seeing deportalization of Facebook. The desktop web interface for Facebook is a portal style interface providing access to all the features made available by the backend. It's big, complicated, and creaky. Who really loves the Facebook UI?
When Facebook moved to mobile they tried the portal approach and it didn't work. So they are going with a strategy of smaller, more focussed, purpose built apps. Mobile first! There's only so much you can do on a small screen. On mobile it's easier to go find a special app than it is to find a menu buried deep within a complicated portal style application.
But Facebook is going one step further. They are not only creating purpose built apps, they are providing multiple competing apps that provide similar functionality and these apps may not even share a backend infrastructure. We see this with Messenger and WhatsApp, Instagram and Facebook's photo app. Paper is an alternate interface to Facebook that provides very limited functionality, but it does what it does very well.
Conway's law may be operating here. The idea that "organizations which design systems ... are constrained to produce designs which are copies of the communication structures of these organizations." With a monolithic backend infrastructure we get a Borg-like portal design. The move to mobile frees the organization from this way of thinking. If apps can be built that provide a view of just a slice of the Facebook infrastructure then apps can be built that don't use Facebook's infrastructure at all. And if they don't need Facebook's infrastructure then they are free not to be built by Facebook at all. So exactly what is Facebook then?
Facebook CEO Mark Zuckerberg has his own take, saying in a keynote presentation at the Mobile World Congress that Facebook's acquisition of WhatsApp was closely related to the Internet.org vision:
The idea is to develop a group of basic internet services that would be free of charge to use — "a 911 for the internet." These could be a social networking service like Facebook, a messaging service, maybe search and other things like weather. Providing a bundle of these free of charge to users will work like a gateway drug of sorts — users who may be able to afford data services and phones these days just don't see the point of why they would pay for those data services. This would give them some context for why they are important, and that will lead them to paying for more services like this — or so the hope goes.
This is the long play, which is a game that having a huge reservoir of valuable stock allows you to play.
Have we reached a conclusion? I don't think so. It's such a stunning dollar amount with such tenuous apparent immediate rewards, that the long term play explanation actually does make some sense. We are still in the very early days of mobile. Nobody knows what the future will look like, so it pays not try to force the future to look like your past. Facebook seems to be doing just that.
But enough of this. How do you support 450 million active users with only 32 engineers? Let's find out...

Sources

A warning here, we don't know a lot about the WhatsApp over all architecture. Just bits and pieces gathered from various sources. Rick Reed's main talk is about the optimization process used to get to 2 million connections a server while using Erlang, which is interesting, but it's not a complete architecture talk.

Stats

These stats are generally for the current system, not the system we have a talk on. The talk on the current system will include more on hacks for data storage, messaging, meta-clustering, and more BEAM/OTP patches.
  • 450 million active users, and reached that number faster than any other company in history.
  • 32 engineers, one developer supports 14 million active users
  • 50 billion messages every day across seven platforms (inbound + outbound)
  • 1+ million people sign up every day
  • $0 invested in advertising
  • $8 million investment
  • Hundreds of nodes
  • >8000 cores
  • Hundreds of terabytes of RAM
  • >70M Erlang messages per second
  • In 2011 WhatsApp achieved 1 million established tcp sessions on a single machine with memory and cpu to spare. In 2012 that was pushed to over 2 million tcp connections. In 2013 WhatsApp tweeted out: On Dec 31st we had a new record day: 7B msgs inbound, 11B msgs outbound = 18 billion total messages processed in one day! Happy 2013!!!

Platform

Backend

  • Erlang
  • FreeBSD
  • Yaws, lighttpd
  • Custom patches to BEAM (BEAM is like Java's JVM, but for Erlang)
  • Custom XMPP

Frontend

  • Seven client platforms: iPhone, Android, Blackberry, Nokia Symbian S60, Nokia S40, Windows Phone, ?
  • SQLite

Hardware

  • Standard user facing server:
    • Dual Westmere Hex-core (24 logical CPUs);
    • 100GB RAM, SSD;
    • Dual NIC (public user-facing network, private back-end/distribution);

Product

  • Focus is on messaging. Connecting people all over the world, regardless of where they are in the world, without having to pay a lot of money. Founder Jan Koum remembers how difficult it was in 1992 to connect to family all over the world.
  • Privacy. Shaped by Jan Koum's experiences growing up in the Ukraine, where nothing was private. Messages are not stored on servers; chat history is not stored; goal is to know as little about users as possible; your name and your gender are not known; chat history is only on your phone.

General

  • WhatsApp server is almost completely implemented in Erlang.
    • Server systems that do the backend message routing are done in Erlang.
    • Great achievement is that the number of active users is managed with a really small server footprint. Team consensus is that it is largely because of Erlang.
    • Interesting to note Facebook Chat was written in Erlang in 2009, but they went away from it because it was hard to find qualified programmers.
  • WhatsApp server has started from ejabberd
    • Ejabberd is a famous open source Jabber server written in Erlang.
    • Originally chosen because its open, had great reviews by developers, ease of start and the promise of Erlang's long term suitability for large communication system.
    • The next few years were spent re-writing and modifying quite a few parts of ejabberd, including switching from XMPP to internally developed protocol, restructuring the code base and redesigning some core components, and making lots of important modifications to Erlang VM to optimize server performance.
  • To handle 50 billion messages a day the focus is on making a reliable system that works. Monetization is something to look at later, it's far far down the road.
  • A primary gauge of system health is message queue length. The message queue length of all the processes on a node is constantly monitored and an alert is sent out if they accumulate backlog beyond a preset threshold. If one or more processes falls behind that is alerted on, which gives a pointer to the next bottleneck to attack.
  • Multimedia messages are sent by uploading the image, audio or video to be sent to an HTTP server and then sending a link to the content along with its Base64 encoded thumbnail (if applicable).
  • Some code is usually pushed every day. Often, it's multiple times a day, though in general peak traffic times are avoided. Erlang helps being aggressive in getting fixes and features into production. Hot-loading means updates can be pushed without restarts or traffic shifting. Mistakes can usually be undone very quickly, again by hot-loading. Systems tend to be much more loosely-coupled which makes it very easy to roll changes out incrementally.
  • What protocol is used in Whatsapp app? SSL socket to the WhatsApp server pools. All messages are queued on the server until the client reconnects to retrieve the messages. The successful retrieval of a message is sent back to the whatsapp server which forwards this status back to the original sender (which will see that as a "checkmark" icon next to the message). Messages are wiped from the server memory as soon as the client has accepted the message
  • How does the registration process work internally in Whatsapp? WhatsApp used to create a username/password based on the phone IMEI number. This was changed recently. WhatsApp now uses a general request from the app to send a unique 5 digit PIN. WhatsApp will then send a SMS to the indicated phone number (this means the WhatsApp client no longer needs to run on the same phone). Based on the pin number the app then request a unique key from WhatsApp. This key is used as "password" for all future calls. (this "permanent" key is stored on the device). This also means that registering a new device will invalidate the key on the old device.
  • Google's push service is used on Android.
  • More users on Android. Android is more enjoyable to work with. Developers are able to prototype a feature and push it out to hundreds of millions of users overnight, if there's an issue it can be fixed quickly. iOS, not so much.

The Quest for 2+ Million Connections Per Server

  • Experienced lots of user growth, which is a good problem to have, but it also means having to spend money buying more hardware and increased operational complexity of managing all those machines.
  • Need to plan for bumps in traffic. Examples are soccer games and earthquakes in Spain or Mexico. These happen near peak traffic loads, so there needs to be enough spare capacity to handle peaks + bumps. A recent soccer match generated a 35% spike in outbound message rate right at the daily peak.
  • Initial server loading was 200 simultaneous connections per server.
    • Extrapolated out would mean a lot of servers with the hoped for growth pattern.
    • Servers were brittle in the face of burst loads. Network glitches and other problems would occur. Needed to decouple components so things weren't so brittle at high capacity.
    • Goal was a million connections per server. An ambitious goal given at the time they were running at 200K connections. Running servers with headroom to allow for world events, hardware failures, and other types of glitches would require enough resilience to handle the high usage levels and failures.

Tools and Techniques Used to Increase Scalability

  • Wrote system activity reporter tool (wsar):
    • Record system stats across the system, including OS stats, hardware stats, BEAM stats. It was build so it was easy to plugin metrics from other systems, like virtual memory. Track CPU utilization, overall utilization, user time, system time, interrupt time, context switches, system calls, traps, packets sent/received, total count of messages in queues across all processes, busy port events, traffic rate, bytes in/out, scheduling stats, garbage collection stats, words collected, etc.
    • Initially ran once a minute. As the systems were driven harder one second polling resolution was required because events that happened in the space if a minute were invisible. Really fine grained stats to see how everything is performing.
  • Hardware performance counters in CPU (pmcstat):
    • See where the CPU is at as a percentage of time. Can tell how much time is being spent executing the emulator loop. In their case it is 16% which tells them that only 16% is executing emulated code so even if you were able to remove all the execution time of all the Erlang code it would only save 16% out of the total runtime. This implies you should focus in other areas to improve efficiency of the system.
  • dtrace, kernel lock-counting, fprof
    • Dtrace was mostly for debugging, not performance.
    • Patched BEAM on FreeBSD to include CPU time stamp.
    • Wrote scripts to create an aggregated view of across all processes to see where routines are spending all the  time.
    • Biggest win was compiling the emulator with lock counting turned on.
  • Some Issues:
    • Earlier on saw more time spent in the garbage collections routines, that was brought down.
    • Saw some issues with the networking stack that was tuned away.
    • Most issues were with lock contention in the emulator which shows strongly in the output of the lock counting.
  • Measurement:
    • Synthetic workloads, which means generating traffic from your own test scripts, is of little value for tuning user facing systems at extreme scale.
      • Worked well for simple interfaces like a user table, generating inserts and reads as quickly as possible.
      • If supporting a million connections on a server it would take 30 hosts to open enough IP ports to generate enough connections to test just one server. For two million servers that would take 60 hosts. Just difficult to generate that kind of scale.
      • The type of traffic that is seen during production is difficult to generate. Can guess at a normal workload, but in actuality see networking events, world events, since multi-platform see varying behaviour between clients, and varying by country.
    • Tee'd workload:
      • Take normal production traffic and pipe it off to a separate system.
      • Very useful for systems for which side effects could be constrained. Don't want to tee traffic and do things that would affect the permanent state of a user or result in multiple messages going to users.
      • Erlang supports hot loading, so could be under a full production load, have an idea, compile, load the change as the program is running and instantly see if that change is better or worse.
      • Added knobs to change production load dynamically and see how it would affect performance. Would be tailing the sar output looking at things like CPU usage, VM utilization, listen queue overflows, and turn knobs to see how the system reacted.
    • True production loads:
      • Ultimate test. Doing both input work and output work.
      • Put server in DNS a couple of times so it would get double or triple the normal traffic. Creates issues with TTLs because clients don't respect DNS TTLs and there's a delay, so can't quickly react to getting more traffic than can be dealt with.
      • IPFW. Forward traffic from one server to another so could give a host exactly the number of desired client connections. A bug caused a kernel panic so that didn't work very well.
  • Results:
    • Started at 200K simultaneous connections per server.
    • First bottleneck showed up at 425K. System ran into a lot of contention. Work stopped. Instrumented the scheduler to measure how much useful work is being done, or sleeping, or spinning. Under load it started to hit sleeping locks so 35-45% CPU was being used across the system but the schedulers are at 95% utilization.
    • The first round of fixes got to over a million connections.
      • VM usage is at 76%. CPU is at 73%. BEAM emulator running at 45% utilization, which matches closely to user percentage, which is good because the emulator runs as user.
      • Ordinarily CPU utilization isn't a good measure of how busy a system is because the scheduler uses CPU.
    • A month later tackling bottlenecks 2 million connections per server was achieved.
      • BEAM utilization at 80%, close to where FreeBSD might start paging. CPU is about the same, with double the connections. Scheduler is hitting contention, but running pretty well.
    • Seemed like a good place to stop so started profiling Erlang code.
      • Originally had two Erlang processes per connection. Cut that to one.
      • Did some things with timers.
    • Peaked at 2.8M connections per server
      • 571k pkts/sec, >200k dist msgs/sec
      • Made some memory optimizations so VM load was down to 70%.
    • Tried 3 million connections, but failed.
      • See long message queues when the system is in trouble. Either a single message queue or a sum of message queues.
      • Added to BEAM instrumentation on message queue stats per process. How many messages are being sent/received, how fast.
      • Sampling every 10 seconds, could see a process had 600K messages in its message queue with a dequeue rate of 40K with a delay of 15 seconds. Projected drain time was 41 seconds.
  • Findings:
    • Erlang + BEAM + their fixes - has awesome SMP scalability. Nearly linear scalability. Remarkable. On a 24-way box can run the system with 85% CPU utilization and it's keeping up running a production load. It can run like this all day.
      • Testament to Erlang's program model.
      • The longer a server has been up it will accumulate long running connections that are mostly idle so it can handle more connections because these connections are not as busy per connection.
    • Contention was biggest issue.
      • Some fixes were in their Erlang code to reduce BEAM's contention issues.
      • Some patched to BEAM.
      • Partitioning workload so work didn't have to cross processors a lot.
      • Time-of-day lock. Every time a message is delivered from a port it looks to update the time-of-day which is a single lock across all schedulers which means all CPUs are hitting one lock.
      • Optimized use of timer wheels. Removed bif timer
      • Check IO time table grows arithmetically. Created VM thrashing has the hash table would be reallocated at various points. Improved to use geometric allocation of the table.
      • Added write file that takes a port that you already have open to reduce port contention.
      • Mseg allocation is single point of contention across all allocators. Make per scheduler.
      • Lots of port transactions when accepting a connection. Set option to reduce expensive port interactions.
      • When message queue backlogs became large garbage collection would destabilize the system. So pause GC until the queues shrunk.
    • Avoiding some common things that come at a price.
      • Backported a TSE time counter from FreeBSD 9 to 8. It's a cheaper to read timer. Fast to get time of day, less expensive than going to a chip.
      • Backported igp network driver from FreeBSD 9 because having issue with multiple queue on NICs locking up.
      • Increase number of files and sockets.
      • Pmcstat showed a lot of time was spent looking up PCBs in the network stack. So bumped up the size of the hash table to make lookups faster.
    • BEAM Patches
      • Previously mentioned instrumentation patches. Instrument scheduler to get utilization information, statistics for message queues, number of sleeps, send rates, message counts, etc. Can be done in Erlang code with procinfo, but with a million connections it's very slow.
      • Stats collection is very efficient to gather so they can be run in production.
      • Stats kept at 3 different decay intervals: 1, 10, 100 second intervals. Allows seeing issues over time.
      • Make lock counting work for larger async thread counts.
      • Added debug options to debug lock counters.
    • Tuning
      • Set the scheduler wake up threshold to low because schedulers would go to sleep and would never wake up.
      • Prefer mseg allocators over malloc.
      • Have an allocator per instance per scheduler.
      • Configure carrier sizes start out big and get bigger. Causes FreeBSD to use super pages. Reduced TLB thrash rate and improves throughput for the same CPU.
      • Run BEAM at real-time priority so that other things like cron jobs don't interrupt schedule. Prevents glitches that would cause backlogs of important user traffic.
      • Patch to dial down spin counts so the scheduler wouldn't spin.
    • Mnesia
      • Prefer os:timestamp to erlang:now.
      • Using no transactions, but with remote replication ran into a backlog. Parallelized replication for each table to increase throughput.
    • There are actually lots more changes that were made.

Lessons

  • Optimization is dark grungy work suitable only for trolls and engineers. When Rick is going through all the changes that he made to get to 2 million connections a server it was mind numbing. Notice the immense amount of work that went into writing tools, running tests, backporting code, adding gobs of instrumentation to nearly every level of the stack, tuning the system, looking at traces, mucking with very low level details and just trying to understand everything. That's what it takes to remove the bottlenecks in order to increase performance and scalability to extreme levels.
  • Get the data you need. Write tools. Patch tools. Add knobs. Ken was relentless in extending the system to get the data they needed, constantly writing tools and scripts to the data they needed to manage and optimize the system. Do whatever it takes.
  • Measure. Remove Bottlenecks. Test. Repeat. That's how you do it.
  • Erlang rocks! Erlang continues to prove its capability as a versatile, reliable, high-performance platform. Though personally all the tuning and patching that was required casts some doubt on this claim.
  • Crack the virality code and profit. Virality is an allusive quality, but as WhatsApp shows, if you do figure out, man, it's worth a lot of money.
  • Value and employee count are now officially divorced. There are a lot of force-multipliers out in the world today. An advanced global telecom infrastructure makes apps like WhatsApp possible. If WhatsApp had to make the network or a phone etc it would never happen. Powerful cheap hardware and Open Source software availability is of course another multiplier. As is being in the right place at the right time with the right product in front of the right buyer.
  • There's something to this brutal focus on the user idea. WhatsApp is focussed on being a simple messaging app, not being a gaming network, or an advertising network, or a disappearing photos network. That worked for them. It guided their no ads stance, their ability to keep the app simple while adding features, and the overall no brainer it just works philosohpy on any phone.
  • Limits in the cause of simplicity are OK. Your identity is tied to the phone number, so if you change your phone number your identity is gone. This is very un-computer like. But it does make the entire system much simpler in design.
  • Age ain't no thing. If it was age discrimination that prevented WhatsApp co-founder Brian Acton from getting a job at both Twitter and Facebook in 2009, then shame, shame, shame.
  • Start simply and then customize. When chat was launched initially the server side was based on ejabberd. It's since been completely rewritten, but that was the initial step in the Erlang direction. The experience with the scalability, reliability, and operability of Erlang in that initial use case led to broader and broader use.
  • Keep server count low. Constantly work to keep server counts as low as possible while leaving enough headroom for events that create short-term spikes in usage. Analyze and optimize until the point of diminishing returns is hit on those efforts and then deploy more hardware.
  • Purposely overprovision hardware. This ensures that users have uninterrupted service during their festivities and employees are able to enjoy the holidays without spending the whole time fixing overload issues.
  • Growth stalls when you charge money. Growth was super fast when WhatsApp was free, 10,000 downloads a day in the early days. Then when switching over to paid that declined to 1,000 a day. At the end of the year, after adding picture messaging, they settled on charging a one-time download fee, later modified to an annual payment.
  • Inspiration comes from the strangest places. Experience with forgetting the username and passwords on Skype accounts drove the passion for making the app "just work."

Related Articles


Evernote helps you remember everything and get organized effortlessly. Download Evernote.

Monday, 9 December 2013

PayPal Switches from Java to JavaScript


PayPal Switches from Java to JavaScript

by Abel Avram on Nov 29, 2013 | Discuss
PayPal has decided to use JavaScript from browser all the way to the back-end server for web applications, giving up legacy code written in JSP/Java.
Jeff Harrell, Director of Engineering at PayPal, has explained in a couple of blog posts (Set My UI Free Part 1: Dust JavaScript Templating, Open Source and More , Node.js at PayPal ) why they decided and some conclusions resulting from switching their web application development from Java/JSP to a complete JavaScript/Node.js stack.
According to Harrell, PayPal's websites had accumulated a good deal of technical debt, and they wanted a "technology stack free of this which would enable greater product agility and innovation." Initially, there was a significant divide between front-end engineers working in web technologies and back-end ones coding in Java. When a UX person wanted to sketch up some pages, they had to ask Java programmers to do some back-end wiring to make it work. This did not fit with their Lean UX development model:
At the time, our UI applications were based on Java and JSP using a proprietary solution that was rigid, tightly coupled and hard to move fast in. Our teams didn't find it complimentary to our Lean UX development model and couldn't move fast in it so they would build their prototypes in a scripting language, test them with users, and then later port the code over to our production stack.
They wanted a "templating [solution that] must be decoupled from the underlying server technology and allow us to evolve our UIs independent of the application language" and that would work with multiple environments. They decided to go with Dust.js  – a templating framework backed up by LinkedIn – , plus Twitter's Bootstrap  and Bower , a package manager for the web. Additional pieces added later were LESS , RequireJS , Backbone.js , Grunt , and Mocha .
Some of PayPal's pages have been redesigned but they still had some of the legacy stack:
… we have legacy C++/XSL and Java/JSP stacks, and we didn't want to leave these UIs behind as we continued to move forward. JavaScript templates are ideal for this. On the C++ stack, we built a library that used V8 to perform Dust renders natively – this was amazingly fast! On the Java side, we integrated Dust using a Spring ViewResolver coupled with Rhino to render the views.
At that time, they also started using Node.js for prototyping new pages, concluding that it was "extremely proficient" and decided to try it in production. For that they also built Kraken.js , a "convention layer" placed on top of Express  which is a Node.js-based web framework. (PayPal has recently open sourced Kraken.js.) The first application to be done in Node.js was the account overview page, which is one of the most accessed PayPal pages, according to Harrell. But because they were afraid the app might not scale well, they decided to create an equivalent Java application to fall back to in case the Node.js one won't work. Following are some conclusions regarding the development effort required for both apps:
Java/Spring JavaScript/Node.js
Set-up time 0 2 months
Development ~5 months ~3 months
Engineers 5 2
Lines of code unspecified 66% of unspecified
The JavaScript team needed 2 months for the initial setup of the infrastructure, but they created with fewer people an application with the same functionality in less time. Running the test suite on production hardware, they concluded that the Node.js app was performing better than the Java one, serving:
Double the requests per second vs. the Java application. This is even more interesting because our initial performance results were using a single core for the node.js application compared to five cores in Java. We expect to increase this divide further.
and having
35% decrease in the average response time for the same page. This resulted in the pages being served 200ms faster— something users will definitely notice.
As a result, PayPal began using the Node.js application in beta in production, and have decided that "all of our consumer facing web applications going forward will be built on Node.js," while some of the existing ones are being ported to Node.js.
One of the benefits of using JavaScript from browser to server is, according to Harrell, the elimination of a divide between front and back-end development by having one team "which allows us to understand and react to our users' needs at any level in the technology stack."

Tell us what you think


Sent from Evernote

Friday, 22 February 2013

Apache Mahout: Scalable machine learning and data mining

Apache Mahout: Scalable machine learning and data mining:


The Apache Mahout™ machine learning library's goal is to build scalable machine learning libraries.

Mahout currently has

  • Collaborative Filtering
  • User and Item based recommenders
  • K-Means, Fuzzy K-Means clustering
  • Mean Shift clustering
  • Dirichlet process clustering
  • Latent Dirichlet Allocation
  • Singular value decomposition
  • Parallel Frequent Pattern mining
  • Complementary Naive Bayes classifier
  • Random forest decision tree based classifier
  • High performance java collections (previously colt collections)
  • A vibrant community
  • and many more cool stuff to come by this summer thanks to Google summer of code
With scalable we mean:
Scalable to reasonably large data sets. Our core algorithms for clustering, classfication and batch based collaborative filtering are implemented on top of Apache Hadoop using the map/reduce paradigm. However we do not restrict contributions to Hadoop based implementations: Contributions that run on a single node or on a non-Hadoop cluster are welcome as well. The core libraries are highly optimized to allow for good performance also for non-distributed algorithms
Scalable to support your business case. Mahout is distributed under a commercially friendly Apache Software license.
Scalable community. The goal of Mahout is to build a vibrant, responsive, diverse community to facilitate discussions not only on the project itself but also on potential use cases. Come to the mailing lists to find out more.
Currently Mahout supports mainly four use cases: Recommendation mining takes users' behavior and from that tries to find items users might like. Clustering takes e.g. text documents and groups them into groups of topically related documents. Classification learns from exisiting categorized documents what documents of a specific category look like and is able to assign unlabelled documents to the (hopefully) correct category. Frequent itemset mining takes a set of item groups (terms in a query session, shopping cart content) and identifies, which individual items usually appear together.

Thursday, 17 May 2012

Article: If all these new DBMS technologies are so scalable, why are Oracle and DB2 still on top of TPC-C? A roadmap to end their dominance.


http://dbmsmusings.blogspot.com/2012/05/if-all-these-new-dbms-technologies-are.html

(This post is coau­thored by Alexan­der Thom­son and Daniel Abadi)
In the last decade, data­base tech­nol­o­gy has arguably pro­gressed fur­thest along the scal­a­bil­i­ty dimen­sion. There have been hun­dreds of research papers, dozens of open-source projects, and numer­ous star­tups attempt­ing to improve the scal­a­bil­i­ty of data­base tech­nol­o­gy. Many of these new tech­nolo­gies have been extreme­ly influential---some papers have earned thou­sands of cita­tions, and some new sys­tems have been deployed by thou­sands of enter­pris­es.

So let's ask a sim­ple ques­tion: If all these new tech­nolo­gies are so scal­able, why on earth are Ora­cle and DB2 still on top of the TPC-C stand­ings? Go to the TPC-C Web­site with the top 10 results in raw trans­ac­tions per sec­ond. As of today (May 16th, 2012), Ora­cle 11g is used for 3 of the results (includ­ing the top result), 10g is used for 2 of the results, and the rest of the top 10 is filled with var­i­ous ver­sions of DB2. How is tech­nol­o­gy designed decades ago still dom­i­nat­ing TPC-C? What hap­pened to all these new tech­nolo­gies with all these scal­a­bil­i­ty claims?

The sur­pris­ing truth is that these new DBMS tech­nolo­gies are not list­ed in theTPC-C top ten results not because that they do not care enough to enter, but rather because they would not win if they did.

To under­stand why this is the case, one must under­stand that scal­a­bil­i­ty does not come for free. Some­thing must be sac­ri­ficed to achieve high scal­a­bil­i­ty. Today, there are three major cat­e­gories of trade­off that can be exploit­ed to make a sys­tem scale. The new tech­nolo­gies basi­cal­ly fall into two of these cat­e­gories; Ora­cle and DB2 fall into a third. And the later parts of this blog post describes research from our group at Yale that intro­duces a fourth cat­e­go­ry of trade­off that pro­vides a roadmap to end the dom­i­nance of Ora­cle and DB2.

These cat­e­gories are:

(1) Sac­ri­fice ACID for scal­a­bil­i­ty. Our pre­vi­ous post on this topic dis­cussed this in detail. Basi­cal­ly we argue that a major class of new scal­able tech­nolo­gies fall under the cat­e­go­ry of "NoSQL" which achieves scal­a­bil­i­ty by drop­ping ACID guar­an­tees, there­by allow­ing them to eschew two phase lock­ing, two phase com­mit, and other imped­i­ments to con­cur­ren­cy and proces­sor inde­pen­dence that hurt scal­a­bil­i­ty. All of these sys­tems that relax ACID are imme­di­ate­ly inel­i­gi­ble to enter the TPC-C com­pe­ti­tion since ACID guar­an­tees are one of TPC-C's require­ments. That's why you don't see NoSQL data­bas­es in the TPC-C top 10---they are imme­di­ate­ly dis­qual­i­fied.

(2) Reduce trans­ac­tion flex­i­bil­i­ty for scal­a­bil­i­ty. There are many so-called"NewSQL" data­bas­es that claim to be both ACID-compliant and scal­able. And these claims are true---to a degree. How­ev­er, the fine print is that they are only lin­ear­ly scal­able when trans­ac­tions can be com­plete­ly iso­lat­ed to a sin­gle "par­ti­tion" or "shard" of data. While these NewSQL data­bas­es often hide the com­plex­i­ty of shard­ing from the appli­ca­tion devel­op­er, they still rely on the shards to be fair­ly inde­pen­dent. As soon as a trans­ac­tion needs to span mul­ti­ple shards (e.g., update two dif­fer­ent user records on two dif­fer­ent shards in the same atom­ic trans­ac­tion), then these NewSQL sys­tems all run into prob­lems. Some sim­ply reject such trans­ac­tions. Oth­ers allow them, but need to per­form two phase com­mit or other agree­ment pro­to­cols in order to ensure ACID com­pli­ance (since each shard may fail inde­pen­dent­ly). Unfor­tu­nate­ly, agree­ment pro­to­cols such as two phase com­mit come at a great scal­a­bil­i­ty cost (see our 2010 paper that explains why). There­fore, NewSQL data­bas­es only scale well if multi-shard trans­ac­tions (also called "dis­trib­uted trans­ac­tions" or "multi-partition trans­ac­tions") are very rare. Unfor­tu­nate­ly for these data­bas­es, TPC-C mod­els a fair­ly rea­son­able retail appli­ca­tion where cus­tomers buy prod­ucts and the inven­to­ry needs to be updat­ed in the same atom­ic trans­ac­tion. 10% of TPC-C New Order trans­ac­tions involve cus­tomers buy­ing prod­ucts from a "remote" ware­house, which is gen­er­al­ly stored in a sep­a­rate shard. There­fore, even for basic appli­ca­tions like TPC-C, NewSQL data­bas­es lose their scal­a­bil­i­ty advan­tages. That's why the NewSQL data­bas­es do not enter TPC-C results --- even just 10% of multi-shard trans­ac­tions caus­es their per­for­mance to degrade rapid­ly.

(3) Trade cost for scal­a­bil­i­ty. If you use high end hard­ware, it is pos­si­ble to get stun­ning­ly high trans­ac­tion­al through­put using old data­base tech­nolo­gies that don't have shared-nothing hor­i­zon­tal­ly scal­a­bil­i­ty. Ora­cle tops TPC-C with an incred­i­bly high through­put of 500,000 trans­ac­tions per sec­ond. There exists no appli­ca­tion in the mod­ern world that pro­duces more than 500,000 trans­ac­tions per sec­ond (as long as humans are ini­ti­at­ing the transactions---machine-generated trans­ac­tions are a dif­fer­ent story). There­fore, Ora­cle basi­cal­ly has all the scal­a­bil­i­ty that is need­ed for human scale appli­ca­tions. The only down­side is cost---the Ora­cle sys­tem that is able to achieve 500,000 trans­ac­tions per sec­ond costs a pro­hib­i­tive $30,000,000!

Since the first two types of trade­offs are imme­di­ate dis­qual­i­fiers for TPC-C, the only remain­ing thing to give up is cost-for-scale, and that's why the old data­base tech­nolo­gies are still dom­i­nat­ing TPC-C. None of these new tech­nolo­gies can han­dle both ACID and 10% remote trans­ac­tions.

A fourth approach...

TPC-C is a very rea­son­able appli­ca­tion. New tech­nolo­gies should be able to han­dle it. There­fore, at Yale we set out to find a new dimen­sion in this trade­off space that could allow a sys­tem to han­dle TPC-C at scale with­out cost­ing $30,000,000. Indeed, we are pre­sent­ing a paper next week at SIG­MOD (see the full paper) that describes a sys­tem that can achieve 500,000 ACID-compliant TPC-C New Order trans­ac­tions per sec­ond using com­mod­i­ty hard­ware in the cloud. The cost to us to run these exper­i­ments was less than $300 (of course, this is rent­ing hard­ware rather than buy­ing, so it's hard to com­pare prices --- but still --- a fac­tor of 100,000 less than $30,000,000 is quite large).

Calvin, our pro­to­type sys­tem designed and built by a large team of researchers at Yale that include Thad­deus Dia­mond, Shu-Chun Weng, Kun Ren, Philip Shao, Anton Petrov, Michael Giuf­fri­da, and Aaron Segal (in addi­tion to the authors of this blog post), explores a trade­off very dif­fer­ent from the three described above. Calvin requires all trans­ac­tions to be exe­cut­ed fully server-side and sac­ri­fices the free­dom to non-deterministically abort or reorder trans­ac­tions on-the-fly dur­ing exe­cu­tion. In return, Calvin gets scal­a­bil­i­ty, ACID-compliance, and extreme­ly low-overhead multi-shard trans­ac­tions over a shared-nothing archi­tec­ture. In other words, Calvin is designed to han­dle high-volume OLTP through­put on shard­ed data­bas­es on cheap, com­mod­i­ty hard­ware stored local­ly or in the cloud. Calvin sig­nif­i­cant­lyimproves the scal­a­bil­i­ty over our pre­vi­ous approach to achiev­ing deter­min­ism in data­base sys­tems.

Scal­ing ACID

The key to Calvin's strong per­for­mance is that it reor­ga­nizes the trans­ac­tion exe­cu­tion pipeline nor­mal­ly used in DBMSs accord­ing to the prin­ci­ple: do all the "hard" work before acquir­ing locks and begin­ning exe­cu­tion. In par­tic­u­lar, Calvin moves the fol­low­ing stages to the front of the pipeline:

  • Repli­ca­tion. In tra­di­tion­al sys­tems, repli­cas agree on each mod­i­fi­ca­tion to data­base state only after some trans­ac­tion has made the change at some "mas­ter" repli­ca. In Calvin, all repli­cas agree in advance on the sequence of trans­ac­tions that they will (deter­min­is­ti­cal­ly) attempt to exe­cute.
  • Agree­ment between par­tic­i­pants in dis­trib­uted trans­ac­tions. Data­base sys­tems tra­di­tion­al­ly use two-phase com­mit (2PC) to han­dle dis­trib­uted trans­ac­tions. In Calvin, every node sees the same glob­al sequence of trans­ac­tion requests, and is able to use this already-agreed-upon infor­ma­tion in place of a com­mit pro­to­col.
  • Disk access­es. In our VLDB 2010 paper, we observed that deter­min­is­tic sys­tems per­formed ter­ri­bly in disk-based envi­ron­ments due to hold­ing locks for the 10ms+ dura­tion of read­ing the need­ed data from disk, since they can­not reorder con­flict­ing trans­ac­tions on the fly. Calvin gets around this set­back by prefetch­ing into mem­o­ry all records that a trans­ac­tion will need dur­ing the repli­ca­tion phase---before locks are even acquired.

As a result, each trans­ac­tion's user-specified logic can be exe­cut­ed at each shard with an absolute min­i­mum of run­time syn­chro­niza­tion between shards or repli­cas to slow it down, even if the trans­ac­tion's logic requires it to access records at mul­ti­ple shards. By min­i­miz­ing the time that locks are held, con­cur­ren­cy can be great­ly increased, there­by lead­ing to near-linear scal­a­bil­i­ty on a com­mod­i­ty clus­ter of machines. 

Strong­ly con­sis­tent glob­al repli­ca­tion 

Calvin's deter­min­is­tic exe­cu­tion seman­tics pro­vide an addi­tion­al ben­e­fit: repli­cat­ing trans­ac­tion­al input is suf­fi­cient to achieve strong­ly con­sis­tent repli­ca­tion. Since repli­cat­ing batch­es of trans­ac­tion requests is extreme­ly inex­pen­sive and hap­pens before the trans­ac­tions acquire locks and begin exe­cut­ing, Calvin's trans­ac­tion­al through­put capac­i­ty does not depend at all on its repli­ca­tion con­fig­u­ra­tion. 

In other words, not only can Calvin can run 500,000 trans­ac­tions per sec­ond on 100 EC2 instances in Ama­zon's US East (Vir­ginia) data cen­ter, it can main­tain strongly-consistent, up-to-date 100-node repli­cas in Ama­zon's Europe (Ire­land) and US West (Cal­i­for­nia) data centers---at no cost to through­put. 

Calvin accom­plish­es this by hav­ing repli­cas per­form the actu­al pro­cess­ing of trans­ac­tions com­plete­ly inde­pen­dent­ly of one anoth­er, main­tain­ing strong con­sis­ten­cy with­out hav­ing to con­stant­ly syn­chro­nize trans­ac­tion results between repli­cas. (Calvin's end-to-end trans­ac­tion laten­cy does depend on mes­sage delays between repli­cas, of course---there is no get­ting around the speed of light.) 

Flex­i­ble data model 

So where does Calvin fall in the OldSQL/NewSQL/NoSQL tri­choto­my? 

Actu­al­ly, nowhere. Calvin is not a data­base sys­tem itself, but rather a trans­ac­tion sched­ul­ing and repli­ca­tion coor­di­na­tion ser­vice. We designed the sys­tem to inte­grate with any data stor­age layer, rela­tion­al or oth­er­wise. Calvin allows user trans­ac­tion code to access the data layer freely, using any data access lan­guage or inter­face sup­port­ed by the under­ly­ing stor­age engine (so long as Calvin can observe which records user trans­ac­tions access). The exper­i­ments pre­sent­ed in the paper use a cus­tom key-value store. More recent­ly, we've hooked Calvin up to Google's Lev­elDB and added sup­port for SQL-based data access with­in trans­ac­tions, build­ing rela­tion­al tables on top of Lev­elDB's effi­cient sorted-string stor­age. 

From an appli­ca­tion devel­op­er's point of view, Calvin's pri­ma­ry lim­i­ta­tion com­pared to other sys­tems is that trans­ac­tions must be exe­cut­ed entire­ly server-side. Calvin has to know in advance what code will be exe­cut­ed for a given trans­ac­tion. Users may pre-define trans­ac­tions direct­ly in C++, or sub­mit arbi­trary Python code snip­pets on-the-fly to be parsed and exe­cut­ed as trans­ac­tions. 

For some appli­ca­tions, this require­ment of com­plete­ly server-side trans­ac­tions might be a dif­fi­cult lim­i­ta­tion. How­ev­er, many appli­ca­tions pre­fer to exe­cute trans­ac­tion code on the data­base serv­er any­way (in the form of stored pro­ce­dures), in order to avoid mul­ti­ple round trip mes­sages between the data­base serv­er and appli­ca­tion serv­er in the mid­dle of a trans­ac­tion. 

If this lim­i­ta­tion is accept­able, Calvin presents a nice alter­na­tive in the trade­off space to achiev­ing high scal­a­bil­i­ty with­out sac­ri­fic­ing ACID or multi-shard trans­ac­tions. Hence, we believe that ourSIG­MOD paper may present a roadmap for over­com­ing the scal­a­bil­i­ty dom­i­nance of the decades-old data­base solu­tions on tra­di­tion­al OLTP work­loads. We look for­ward to debat­ing the mer­its of this approach in the weeks ahead (and Alex will be pre­sent­ing the paper at SIG­MOD next week).

Thursday, 19 April 2012

Building Highly Available Systems in Erlang

InfoQ: Building Highly Available Systems in Erlang:

Key ideas:


The process approach to fault isolation advocates that the process
software be fail-fast, it should either function correctly or it
should detect the fault, signal failure and stop operating.

  Processes are  made fail-fast  by defensive programming.  They check
all their inputs, intermediate results and data structures as a matter
of course. If any error is detected, they signal a failure and stop. In
the  terminology of  [Christian],  fail-fast software  has small  fault
detection latency.

Saturday, 7 April 2012

Are Cloud Based Memory Architectures the Next Big Thing?

Are Cloud Based Memory Architectures the Next Big Thing?:
We are on the edge of two potent technological changes: Clouds and Memory Based Architectures. This evolution will rip open a chasm where new players can enter and prosper. Google is the master of disk. You can't beat them at a game they perfected. Disk based databases like SimpleDB and BigTable are complicated beasts, typical last gasp products of any aging technology before a change. The next era is the age of Memory and Cloud which will allow for new players to succeed. The tipping point will be soon.

Let's take a short trip down web architecture lane:

  • It's 1993: Yahoo runs on FreeBSD, Apache, Perl scripts and a SQL database
  • It's 1995: Scale-up the database.
  • It's 1998: LAMP
  • It's 1999: Stateless + Load Balanced + Database + SAN
  • It's 2001: In-memory data-grid.
  • It's 2003: Add a caching layer.
  • It's 2004: Add scale-out and partitioning.
  • It's 2005: Add asynchronous job scheduling and maybe a distributed file system.
  • It's 2007: Move it all into the cloud.
  • It's 2008: Cloud + web scalable database.
  • It's 20??: Cloud + Memory Based Architectures

    You may disagree with the timing of various innovations and you would be correct. I couldn't find a history of the evolution of website architectures, so I just made stuff up. If you have any better information please let me know.

    Why might cloud based memory architectures be the next big thing? For now we'll just address the memory based architecture part of the question, the cloud component is covered a little later.

    Behold the power of keeping data in memory:
    Google query results are now served in under an astonishingly fast 200ms, down from 1000ms in the olden days. The vast majority of this great performance improvement is due to holding indexes completely in memory. Thousands of machines process each query in order to make search results appear nearly instantaneously.
    This text was adapted from notes on Google Fellow Jeff Dean keynote speech at WSDM 2009.

    Google isn't the only one getting a performance bang from moving data into memory. Both LinkedInand Digg keep the graph of their network social network in memory. Facebook has northwards of 800 memcached servers creating a reservoir of 28 terabytes of memory enabling a 99% cache hit rate. Even little guys can handle 100s of millions of events per day by using memory instead of disk.

    With their new Unified Computing strategy Cisco is also entering the memory game. Their new machines "will be focusing on networking and memory" with servers crammed with 384 GB of RAM, fast processors, and blazingly fast processor interconnects. Just what you need when creating memory based systems.

    Memory Is The System Of Record

    What makes Memory Based Architectures different from traditional architectures is that memory is the system of record. Typically disk based databases have been the system of record. Disk has been King, safely storing data away within its castle walls. Disk being slow we've ended up wrapping disks in complicated caching and distributed file systems to make them perform.

    Sure, memory is used as all over the place as cache, but we're always supposed to pretend that cache can be invalidated at any time and old Mr. Reliable, the database, will step in and provide the correct values. In Memory Based Architectures memory is where the "official" data values are stored.

    Caching also serves a different purpose. The purpose behind cache based architectures is to minimize the data bottleneck through to disk. Memory based architectures can address the entire end-to-end application stack. Data in memory can be of higher reliability and availability than traditional architectures.

    Memory Based Architectures initially developed out of the need in some applications spaces for very low latencies. The dramatic drop of RAM prices along with the ability of servers to handle larger and larger amounts of RAM has caused memory architectures to verge on going mainstream. For example, someone recently calculated that 1TB of RAM across 40 servers at 24 GB per server would cost an additional $40,000. Which is really quite affordable given the cost of the servers. Projecting out, 1U and 2U rack-mounted servers will soon support a terabyte or more or memory.

    RAM = High Bandwidth And Low Latency

    Why are Memory Based Architectures so attractive? Compared to disk RAM is a high bandwidth and low latency storage medium. Depending on who you ask the bandwidth of RAM is 5 GB/s. The bandwidth of disk is about 100 MB/s. RAM bandwidth is many hundreds of times faster. RAM wins. Modern hard drives have latencies under 13 milliseconds. When many applications are queued for disk reads latencies can easily be in the many second range. Memory latency is in the 5 nanosecond range. Memory latency is 2,000 times faster. RAM wins again.

    RAM Is The New Disk

    The superiority of RAM is at the heart of the RAM is the New Disk paradigm. As an architecture it combines the holy quadrinity of computing:
  • Performance is better because data is accessed from memory instead of through a database to a disk.
  • Scalability is linear because as more servers are added data is transparently load balanced across the servers so there is an automated in-memory sharding.
  • Availability is higher because multiple copies of data are kept in memory and the entire system reroutes on failure.
  • Application development is faster because there’s only one layer of software to deal with, the cache, and its API is simple. All the complexity is hidden from the programmer which means all a developer has to do is get and put data.

    Access disk on the critical path of any transaction limits both throughput and latency. Committing a transaction over the network in-memory is faster than writing through to disk. Reading data from memory is also faster than reading data from disk. So the idea is to skip disk, except perhaps as an asynchronous write-behind option, archival storage, and for large files.

    Or Is Disk Is The The New RAM

    To be fair there is also a Disk is the the new RAM, RAM is the New Cache paradigm too. This somewhat counter intuitive notion is that a cluster of about 50 disks has the same bandwidth of RAM, so the bandwidth problem is taken care of by adding more disks.

    The latency problem is handled by reorganizing data structures and low level algorithms. It's as simple as avoiding piecemeal reads and organizing algorithms around moving data to and from memory in very large batches and writing highly parallelized programs. While I have no doubt this approach can be made to work by very clever people in many domains, a large chunk of applications are more time in the random access domain space for which RAM based architectures are a better fit.

    Grids And A Few Other Definitions

    There's a constellation of different concepts centered around Memory Based Architectures that we'll need to understand before we can understand the different products in this space. They include:
  • Compute Grid - parallel execution. A Compute Grid is a set of CPUs on which calculations/jobs/work is run. Problems are broken up into smaller tasks and spread across nodes in the grid. The result is calculated faster because it is happening in parallel.
  • Data Grid - a system that deals with data — the controlled sharing and management of large amounts of distributed data.
  • In-Memory Data Grid (IMDG) - parallel in-memory data storage. Data Grids are scaled horizontally, that is by adding more nodes. Data contention is removed removed by partitioning data across nodes.
  • Colocation - Business logic and object state are colocated within the same process. Methods are invoked by routing to the object and having the object execute the method on the node it was mapped to. Latency is low because object state is not sent across the wire.
  • Grid Computing - Compute Grids + Data Grids
  • Cloud Computing - datacenter + API. The API allows the set of CPUs in the grid to be dynamically allocated and deallocated.

    Who Are The Major Players In This Space?

    With that bit of background behind us, there are several major players in this space (in alphabetical order):
  • Coherence - is a peer-to-peer, clustered, in-memory data management system. Coherence is a good match for applications that need write-behind functionality when working with a database and you require multiple applications have ACID transactions on the database. Java, JavaEE, C++, and .NET.
  • GemFire - an in-memory data caching solution that provides low-latency and near-zero downtime along with horizontal & global scalability. C++, Java and .NET.
  • GigaSpaces - GigaSpaces attacks the whole stack: Compute Grid, Data Grid, Message, Colocation, and Application Server capabilities. This makes for greater complexity, but it means there's less plumbing that needs to be written and developers can concentrate on writing business logic. Java, C, or .Net.
  • GridGain - A compute grid that can operate over many data grids. It specializes in the transparent and low configuration implementation of features. Java only.
  • Terracotta - Terracotta is network-attached memory that allows you share memory and do anything across a cluster. Terracotta works its magic at the JVM level and provides: high availability, an end of messaging, distributed caching, a single JVM image. Java only.
  • WebSphere eXtreme Scale. Operates as an in-memory data grid that dynamically caches, partitions, replicates, and manages application data and business logic across multiple servers.

    This class of products has generally been called In-Memory Data Grids (IDMG), though not all the products fit snugly in this category. There's quite a range of different features amongst the different products.

    I tossed IDMG the acronym in favor of Memory Based Architectures because the "in-memory" part seems redundant, the grid part has given way to the cloud, the "data" part really can include both data and code. And there are other architectures that will exploit memory yet won't be classic IDMG. So I just used Memory Based Architecture as that's the part that counts.

    Given the wide differences between the products there's no canonical architecture. As an example here's a diagram of how GigaSpaces In-Memory-Data-Grid on the Cloud works.





    Some key points to note are:
  • A POJO (Plain Old Java Object) is written through a proxy using a hash-based data routing mechanism to be stored in a partition on a Processing Unit. Attributes of the object are used as a key. This is straightforward hash based partitioning like you would use with memcached.
  • You are operating through GigaSpace's framework/container so they can automatically handle things like messaging, sending change events, replication, failover, master-worker pattern, map-reduce, transactions, parallel processing, parallel query processing, and write-behind to databases.
  • Scaling is accomplished by dividing your objects into more partitions and assigning the partitions to Processing Unit instances which run on nodes-- a scale-out strategy. Objects are kept in RAM and the objects contain both state and behavior. A Service Grid component supports the dynamic creation and termination of Processing Units.

    Not conceptually difficult and familiar to anyone who has used caching systems like memcached. Only is this case memory is not just a cache, it's the system of record.

    Obviously there are a million more juicy details at play, but that's the gist of it. Admittedly GigaSpaces is on the full featured side of the product equation, but from a memory based architecture perspective the ideas should generalize. When you shard a database, for example, you generally lose the ability to execute queries, you have to do all the assembly yourself. By using GigaSpaces framework you get a lot of very high-end features like parallel query processing for free.

    The power of this approach certainly comes in part from familiar concepts like partitioning. But the speed of memory versus disk also allows entire new levels of performance and reliability in a relatively simple and easy to understand and deploy package.

    NimbusDB - The Database In The Cloud

    Jim Starkey, President of NimbusDB, is not following the IDMG gang's lead. He's taking a completely fresh approach based on thinking of the cloud as a new platform unto itself. Starting from scratch, what would a database for the cloud look like?

    Jim is in position to answer this question as he has created a transactional database engine for MySQL named Falcon and added multi-versioning support to InterBase, the first relational database to feature MVCC (Multiversion Concurrency Control).

    What defines the cloud as a platform? Here's are some thoughts from Jim I copied out of the Cloud Computing group. You'll notice I've quoted Jim way way too much. I did that because Jim is an insightful guy, he has a lot of interesting things to say, and I think he has a different spin on the future of databases in the cloud than anyone else I've read. He also has the advantage of course of not having a shipping product, but we shall see.
  • I've probably said this before, but the cloud is a new computing platform that some have learned to exploit, others are scrambling to master, but most people will see as nothing but a minor variation on what they're already doing. This is not new. When time sharing as invented, the batch guys considered it as remote job entry, just a variation on batch. When departmental computing came along (VAXes, et al), the timesharing guys considered it nothing but timesharing on a smaller scale. When PCs and client/server computing came along, the departmental computing guys (i.e. DEC), considered PCs to be a special case of smart terminals. And when the Internet blew into town, the client server guys considered it as nothing more than a global scale LAN. So the batchguys are dead, the timesharing guys are dead, the departmental computing guys are dead, and the client server guys are dead. Notice a pattern?
  • The reason that databases are important to cloud computing is that virtually all applications involve the interaction of client data with a shared, persistent data store. And while application processing can be easily scaled, the limiting factor is the database system. So if you plan to do anything more than play Tetris in the cloud, the issue of database management should be foremost in your mind.
  • Disks are the limiting factors in contemporary database systems. Horrible things, disk. But conventional wisdom is that you build a clustered database system by starting with a distributed file system. Wrong. Evolution is faster processors, bigger memory, better tools. Revolution
    is a different way of thinking, a different topology, a different way of putting the parts together.
  • What I'm arguing is that a cloud is a different platform, and what works well for a single computer doesn't work at all well in cloud, and things that work well in a cloud don't work at all on the single computer system. So it behooves us to re-examine a lot an ancient and honorable assumptions to see if they make any sense at all in this brave new world.
  • Sharing a high performance disk system is fine on a single computer, troublesome in a cluster, and miserable on a cloud.
  • I'm a database guy who's had it with disks. Didn't much like the IBM 1301, and disks haven't gotten much better since. Ugly, warty, slow, things that require complex subsystems to hide their miserable characteristics. The alternative is to use the memory in a cloud as a distributed L2
    cache. Yes, disks are still there, but they're out of the performance loop except for data so stale that nobody has it memory.
  • Another machine or set of machines is just as good as a disk. You can quibble about reliable power, etc, but write queuing disks have the same problem.
  • Once you give up the idea of logs and page caches in favor of asynchronous replications, life gets a great deal brighter. It really does make sense to design to the strengths of cloud(redundancy) rather than their weaknesses (shared anything).
  • And while one guys is fetching his 100 MB per second, the disk is busy and everyone else is waiting in line contemplating existence. Even the cheapest of servers have two gigabit ethernet channels and switch. The network serves everyone in parallel while the disk is single threaded
  • I favor data sharing through a formal abstraction like a relational database. Shared objects are things most programmers are good at handling. The fewer the things that application developers need to manage the more likely it is that the application will work.
  • I buy the model of object level replication, but only as a substrate for something with a more civilized API. Or in other words, it's a foundation, not a house.
  • I'd much rather have a pair of quad-core processors running as independent servers than contending for memory on a dual socket server. I don't object to more cores per processor chip, but I don't want to pay for die size for cores perpetually stalled for memory.
  • The object substrate worries about data distribution and who should see what. It doesn't even know it's a database. SQL semantics are applied by an engine layered on the object substrate. The SQL engine doesn't worry or even know that it's part of a distributed database -- it just executes SQL statements. The black magic is MVCC.
  • I'm a database developing building a database system for clouds. Tell me what you need. Here is my first approximation: A database that scales by adding more computers and degrades gracefully when machines are yanked out; A database system that never needs to be shut down; Hardware and software fault tolerance; Multi-site archiving for disaster survival; A facility to reach into the past to recover from human errors (drop table customers; oops;); Automatic load balancing
  • MySQL scales with read replication which requires a full database copy to start up. For any cloud relevant application, that's probably hundreds of gigabytes. That makes it a mighty poor candidate for on-demand virtual servers.
  • Do remember that the primary function of a database system is to maintain consistency. You don't want a dozen people each draining the last thousand buckets from a bank account or a debit to happen without the corresponding credit.
  • Whether the data moves to the work or the work moves to the data isn't that important as long as they both end up a the same place with as few intermediate round trips as possible.
  • In my area, for example, databases are either limited by the biggest, ugliest machine you can afford *or* you have to learn to operation without consistent, atomic transactions. A bad rock / hard place choice that send the cost of scalable application development through the ceiling. Once we solve that, applications that server 20,000,000 users will be simple and cheap to write. Who knows where that will go?
  • To paraphrase our new president, we must reject the false choice between data consistency and scalability.
  • Cloud computing is about using many computers to scale problems that were once limited by the capabilities of a single computer. That's what makes clouds exciting, at least to me. But most will argue that cloud computing is a better economic model for running many instances of a
    single computer. Bah, I say, bah!
  • Cloud computing is a wonder new platform. Let's not let the dinosaurs waiting for extinction define it as a minor variation of what they've been doing for years. They will, of course, but this (and the dinosaurs) will pass.
  • The revolutionary idea is that applications don't run on a single computer but an elastic cloud of computers that grows and contracts by demand. This, in turn, requires an applications infrastructure that can a) run a single application across as many machines as necessary, and b) run many applications on the same machines without any of the cross talk and software maintenance problems of years past. No, the software infrastructure required to enable this is not mature and certainly not off the shelf, but many smart folks are working on it.
  • There's nothing limiting in relational except the companies that build them. A relational database can scale as well as BigTable and SimpleDB but still be transactional. And, unlike BigTable and SimpleDB, a relational database can model relationships and do exotic things like transferring money from one account to another without "breaking the bank.". It is true that existing relational database systems are largely constrained to single cpu or cluster with a shared file system, but we'll get over that.
  • Personally, I don't like masters any more than I like slaves. I strongly favor peer to peer architectures with no single point of failure. I also believe that database federation is a work-around
    rather than a feature. If a database system had sufficient capacity, reliability, and availability, nobody would ever partition or shard data. (If one database instance is a headache, a million tiny ones is a horrible, horrible migraine.)
  • Logic does need to be pushed to the data, which is why relational database systems destroyed hierarchical (IMS), network (CODASYL), and OODBMS. But there is a constant need to push semantics higher to further reduce the number of round trips between application semantics and the database systems. As for I/O, a database system that can use the cloud as an L2 cache breaks free from dependencies on file systems. This means that bandwidth and cycles are the limiting factors, not I/O capacity.
  • What we should be talking about is trans-server application architecture, trans-server application platforms, both, or whether one will make the other unnecessary.
  • If you scale, you don't/can't worry about server reliability. Money spent on (alleged) server reliability is money wasted.
  • If you view the cloud as a new model for scalable applications, it is a radical change in computing platform. Most people see the cloud through the lens of EC2, which is just another way to run a server that you have to manage and control, then the cloud is little more than a rather
    boring business model. When clouds evolve to point that applications and databases can utilize whatever resources then need to meet demand without the constraint of single machine limitations, we'll have something really neat.
  • On MVCC: Forget about the concept of master. Synchronizing slaves to a master is hopeless. Instead, think of a transaction as a temporal view of database state; different transactions
    will have different views. Certain critical operations must be serialized, but that still doesn't require that all nodes have identical views of database state.
  • Low latency is definitely good, but I'm designing the system to support geographically separated sub-clouds. How well that works under heavy load is probably application specific. If the amount of volatile data common to the sub-clouds is relatively low, it should work just fine provided there is enough bandwidth to handle the replication messages.
  • MVCC tracks multiple versions to provide a transaction with a view of the database consistent with the instant it started while preventing a transaction from updating a piece of data that it could not see. MVCC is consistent, but it is not serializable. Opinions vary between academia and the real world, but most database practitioners recognize that the consistency provided by MVCC is sufficient for programmers of modest skills to product robust applications.
  • MVCC, heretofore, has been limited to single node databases. Applied to the cloud with suitable bookkeeping to control visibility of updates on individual nodes, MVCC is as close to black magic as you are likely to see in your lifetime, enabling concurrency and consistency with mostly non-blocking, asynchronous messaging. It does, however, dispense with the idea that a cloud has at any given point of time a single definitive state. Serializability implemented with record locking is an attempt to make distributed system march in lock-step so that the result is as if there there no parallelism between nodes. MVCC recognizes that parallelism is the key to scalability. Data that is a few microseconds old is not a problem as long as updates don't collide.

    Jim certainly isn't shy with his opinions :-)

    My summary of what he wants to do with NimbusDB is:
  • Make a scalable relational database in the cloud where you can use normal everyday SQL to perform summary functions, define referential integrity, and all that other good stuff.
  • Transactions scale using a distributed version of MVCC, which I do not believe has been done before. This is the key part of the plan and a lot depends on it working.
  • The database is stored primarily in RAM which makes cloud level scaling of an RDBMS possible.
  • The database will handle all the details of scaling in the cloud. To the developer it will look like just a very large highly available database.

    I'm not sure if NimbusDB will support a compute grid and map-reduce type functionality. The low latency argument for data and code collocation is a good one, so I hope it integrates some sort of extension mechanism.

    Why might NimbusDB be a good idea?
  • Keeps simple things simple. Web scale databases like BigTable and SimpleDB make simple things difficult. They are full of quotas, limits, and restrictions because by their very nature they are just a key-value layer on top of a distributed file system. The database knows as little about the data as possible. If you want to build a sequence number for a comment system, for example, it takes complicated sharding logic to remove write contention. Developers are used to SQL and are comfortable working within the transaction model, so the transition to cloud computing would be that much easier. Now, to be fair, who knows if NimbusDB will be able to scale under high load either, but we need to make simple things simple again.
  • Language independence. Notice the that IDMG products are all language specific. They support some combination of .Net/Java/C/C++. This is because they need low level object knowledge to transparently implement their magic. This isn't bad, but it does mean if you use Python, Erlang, Ruby, or any other unsupported language then you are out of luck. As many problems as SQL has, one of its great gifts is programmatic universal access.
  • Separates data from code. Data is forever, code changes all the time. That's one of the common reasons for preferring a database instead of an objectbase. This also dovetails with the language independence issue. Any application can access data from any language and any platform from now and into the future. That's a good quality to have.

    The smart money has been that cloud level scaling requires abandoning relational databases and distributed transactions. That's why we've seen an epidemic of key-value databases and eventually consistent semantics. It will be fascinating to see if Jim's combination of Cloud + Memory + MVCC can prove the insiders wrong.

    Are Cloud Based Memory Architectures The Next Big Thing?

    We've gone through a couple of different approaches to deploying Memory Based Architectures. So are they the next big thing?

    Adoption has been slow because it's new and different and that inertia takes a while to overcome. Historically tools haven't made it easy for early adopters to make the big switch, but that is changing with easier to deploy cloud based systems. And current architectures, with a lot of elbow grease, have generally been good enough.

    But we are seeing a wide convergence on caching as way to make slow disks perform. Truly enormous amounts of effort are going into adding cache and then trying to keep the database and applications all in-sync with cache as bottom up and top down driven changes flow through the system.

    After all that work it's a simple step to wonder why that extra layer is needed when the data could have just as well be kept in memory from the start. Now add the ease of cloud deployments and the ease of creating scalable, low latency applications that are still easy to program, manage, and deploy. Building multiple complicated layers of application code just to make the disk happy will make less and less sense over time.

    We are on the edge of two potent technological changes: Clouds and Memory Based Architectures. This evolution will rip open a chasm where new players can enter and prosper. Google is the master of disk. You can't beat them at a game they perfected. Disk based databases like SimpleDB and BigTable are complicated beasts, typical last gasp products of any aging technology before a change. The next era is the age of Memory and Cloud which will allow for new players to succeed. The tipping point is soon.

    Related Articles

  • GridGain: One Compute Grid, Many Data Grids
  • GridGain vs Hadoop
  • Cameron Purdy: Defining a Data Grid
  • Compute Grids vs. Data Grids
  • Performance killer: Disk I/O by Nathanael Jones
  • RAM is the new disk... by Steven Robbins
  • Talk on disk as the new RAM by Greg Linden
  • Disk-Based Parallel Computation, Rubik's Cube, and Checkpointing by Gene Cooperman, Northeastern Professor, High Performance Computing Lab - Disk is the the new RAM and RAM is the new cache
  • Disk is the new disk by David Hilley.
  • Latency lags bandwidth by David A. Patterson
  • InfoQ Article - RAM is the new disk... by Nati Shalom
  • Tape is Dead Disk is Tape Flash is Disk RAM Locality is King by Jim Gray
  • Product: ScaleOut StateServer is Memcached on Steroids
  • Cameron Purdy: Defining a Data Grid
  • Compute Grids vs. Data Grids
  • Latency is Everywhere and it Costs You Sales - How to Crush it
  • Virtualization for High Performance Computing by Shai Fultheim
  • Multi-Multicore Single System Image / Cloud Computing. A Good Idea? (part 1) by Greg Pfister
  • How do you design and handle peak load on the Cloud ? by Cloudiquity.
  • Defining a Data Grid by Cameron Purdy
  • The Share-Nothing Architecture by Zef Hemel.
  • Scaling memcached at Facebook
  • Cache-aside, write-behind, magic and why it sucks being an Oracle customer by Stefan Norberg.
  • Introduction to Terracotta by Mike
  • The five-minute rule twenty years later, and how flash memory changes the rules by Goetz Graefe