Original URL: http://www.theregister.co.uk/2011/06/30/yahoo_hadoop_and_realtime/

Yahoo! tech boss gazes beyond Hadoop

MapReduce goes only so fast

By Cade Metz

Posted in Business, 30th June 2011 05:16 GMT

Hadoop underpins everything from Facebook to eBay to Yahoo!, but it's not designed for online applications.

Based on Google MapReduce – the distributed software platform that was long used to build the search giant's web index – Hadoop was designed for off-line data-crunching. It's used to feed online services – to, say, build a web index can that can then accessed by a search engine – but it's not suited to driving the sort of "real-time" interactions that happen on the web.

"With the paths that go through Hadoop [at Yahoo!], the latency is about fifteen minutes," Yahoo! chief technology officer Raymie Stata told The Register at the company's annual Hadoop Summit on Wednesday in Santa Clara, California. "In some cases, that's plenty 'real-time'. And we have been trying to get that down to less than fifteen minutes. But it will never be true real-time. It will never be what we call 'next click', where I click and by the time the page loads, the semantic implication of my decision is reflected in the page."

HBase – the open source distributed database that works in tandem with the Hadoop File Systems (HDFS) – is suited to some real-time web applications. But Stata questions whether HBase will ultimately reinvent the real-time web in the way Hadoop has overhauled batch processing. Yahoo! actually uses HBase to drive a service that customizes its homepage for each visitor. But the company is exploring other real-time platforms, including two Yahoo!-driven open source projects: MapReduce Online and S4.

"We have a number of internal solutions to the 'next-click' problem," Stata says. "And there's active dialogue on how to consolidate this, but we don't quite know enough yet to say 'this is the answer'."

A joint project between Yahoo! Research and the University of California, Berkeley, MapReduce Online is essentially an effort to modify the Hadoop framework for use with real-time tasks. "MapReduce is a sequential process. You map all the data. You reduce all the data. You can get the batch size smaller and smaller. But when you get to a batch size of one, it doesn't make sense," Stata says.

"MapReduce Online moves to a windowing model, like windowed SQL, where you map a 'window' of data...You take a window of stuff that comes in during a given time period, and map that and reduce that. The data is always streaming through."

It's unclear whether the platform is actually used on Yahoo!'s live infrastructure. But Stata says that, at least in theory, it can run the same code as runs on Hadoop MapReduce. The problem, he adds, is that the platform is still burdened by the inherent limitations of the MapReduce framework. In some ways, it's still a batch platform.

S4 is similar to MapReduce Online – except it parts ways with the Hadoop framework. "It's a streaming analytics library," says Todd Papaioannou, vice president of cloud architecture at Yahoo!. "You have data coming through along a pipe and you're doing analytics on a window of data as it goes through." S4 now drives the delivery of Yahoo!'s sponsored search results.

Off HBase

Though he's still unsure about the longterm prospects of MapReduce Online and S4, Stata seems particularly cold on the real-time future of HBase, which is based on Google's BigTable. "HBase can be used in a real-time infrastructure. It's not [necessarily suited] to real-time. Where do you host the application code itself? What do you do about fault tolerance?" he says. "For applications where you truly need to be fault tolerant, HBase is not a great solution. It's not transactional. It doesn't have a lot of features that help write robust code in a simple fashion. It puts a lot of burden on the app developer."

The same limitations apply to BigTable, but Google is looking to push through them. The company no longer builds its web index with MapReduce. It uses a new proprietary platform called Percolator that runs atop BigTable, letting the company continually update the index without having to reprocess the entire thing from scratch, and there's a compute framework that lets engineers execute code atop the platform.

Raymie Stata

"MapReduce and other batch-processing systems cannot process small updates individually as they rely on creating large batches for efficiency," reads Google's paper describing Percolator. "By replacing a batch-based indexing system with an indexing system based on incremental processing using Percolator, we process the same number of documents per day, while reducing the average age of documents in Google search results by 50 per cent."

Stata declined to discuss Google's paper. But his take on HBase indicates this is another area where Yahoo! will not follow in Google's footsteps. With Hadoop – a project Yahoo! bootstrapped after hiring project founder Doug Cutting – the company mimicked the Google way, but only in part. Yahoo! chose to open source the platform. Just this week, it renewed its commitment to the open-source project with the creation Hortonworks, a Hadoop startup built around the company's 25 core Hadoop developers. And both MapReduce Online and S4 are open source. Google's platforms are decidedly closed.

Stata believes that web giants like Yahoo! will better serve themselves by exchanging their infrastructure ideas with the open source world, whether those ideas involve distributed batch processing or real-time calculations. "Yahoo! is often a few years ahead, just because of our scale and the complexities we have to deal with, but that's not where we want to gain longterm proprietary value," Stata says. "In fact, if what we do becomes what everyone else does, that actually benefits us." How very un-Google. ®