Yahoo! tech boss gazes beyond Hadoop

MapReduce goes only so fast

Bridging the IT gap between rising business demands and ageing tools

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. ®

Build a business case: developing custom apps

More from The Register

next story
BBC goes offline in MASSIVE COCKUP: Stephen Fry partly muzzled
Auntie tight-lipped as major outage rolls on
iPad? More like iFAD: We reveal why Apple ran off to IBM
But never fear fanbois, you're still lapping up iPhones, Macs
Nadella: Apps must run on ALL WINDOWS – PCs, slabs and mobes
Phone egg, meet desktop chicken - your mother
HP, Microsoft prove it again: Big Business doesn't create jobs
SMEs get lip service - what they need is dinner at the Club
ITC: Seagate and LSI can infringe Realtek patents because Realtek isn't in the US
Land of the (get off scot) free, when it's a foreign owner
Samsung threatens to cut ties with supplier over child labour allegations
Vows to uphold 'zero tolerance' policy on underage workers
Dude, you're getting a Dell – with BITCOIN: IT giant slurps cryptocash
1. Buy PC with Bitcoin. 2. Mine more coins. 3. Goto step 1
There's NOTHING on TV in Europe – American video DOMINATES
Even France's mega subsidies don't stop US content onslaught
You! Pirate! Stop pirating, or we shall admonish you politely. Repeatedly, if necessary
And we shall go about telling people you smell. No, not really
prev story


Seven Steps to Software Security
Seven practical steps you can begin to take today to secure your applications and prevent the damages a successful cyber-attack can cause.
Consolidation: The Foundation for IT Business Transformation
In this whitepaper learn how effective consolidation of IT and business resources can enable multiple, meaningful business benefits.
Designing a Defense for Mobile Applications
Learn about the various considerations for defending mobile applications - from the application architecture itself to the myriad testing technologies.
Build a business case: developing custom apps
Learn how to maximize the value of custom applications by accelerating and simplifying their development.
Consolidation: the foundation for IT and business transformation
In this whitepaper learn how effective consolidation of IT and business resources can enable multiple, meaningful business benefits.