Posts mit dem Label Skalability werden angezeigt. Alle Posts anzeigen
Posts mit dem Label Skalability werden angezeigt. Alle Posts anzeigen

17. August 2008

Big and Bigger Systems

I started to compile a list of guidelines for developing big systems. So, what is a big system?

  • "Big" means 5k - 100k concurrent users.
Of course, there are even bigger systems. There is also
  • "Massive" 50k up to 1 Mio. concurrent and
  • "Mainstream" 500k up to 10 Mio. concurrent and beyond.
The guidelines are for the smallest one the of 3 noteworty categories: "Big". This is the category I am experienced with. 10 or 20k concurrent is significantly different from normal/sparetime/hobby systems. But it is still conventional technology (e.g. on standard LAMP). I can imagine what needs to be done to get a massive system up and keep it running. ButI have never done it, yet. I can only dream of a mainstream scale system, especially a unique and not sharded one (see below).

Examples:
  • "Big": Xing, major online newspapers, Wikipedia, many browser games,
  • "Massive": Second Life, EVE-Online (barely), Runescape (I guess), Habbo, Facebook,
  • "Mainstream": WoW (barely), major IM networks (QQ, ICQ, MSN), Google search, mobile networks.
In other words:
  • "Big" ranges from "substantially over the top of your single server" to "the complete cluster would appear in the Top 500 Supercomputers list". EVE-Online is at the top of "Big" by sheer numbers, but it is on the brink to "Massive" because of its complexity. Xing: more than 10k click around on the web site at any time, sometimes twice as much.
  • "Massive" is really large. Think about thousands or ten times more servers. That's a decent farm. Second Life is at the lower end. Their server numbers are a bit high for the user count, because of their architecture. Facebook is somwhere in the middle with 50 TB cache memory alone. This number is already 20,000 times more than your PC and growing. We are talking about the large web services, companies with 10 B market capitalization. These are "Massive". The big guys.
  • "Mainstream" is unthinkable. Can you imagine half a million servers? If you stack them, they scratch the International Space Station. If you plan a system like this for everyone, then better avoid the hazzle and make it distributed like the Web. It's made for 10% of the world population (registered users not concurrent). There are few such systems on the planet. Few people managed to build in this order of magnitude. I am sure they grew into it with the system. Nothing that you can buy.
Of course there are differences in complexity between a static web site and an MMOG. A static web site, e.g. a newspaper may need much less resources for 10k concurrent readers, than a vitual world for 10k concurrent players. Comunication in an IM network is very different from a MMOG. MMOGs enter the "Massive" category earlier than browser games.

But there are also factors, which make the complexity similar or at least in the same order of magnitude.
  • As an example, 10k HTTP/HTML readers make more network connections (easily 10k conns per sec.) than 10k players (100 conns per sec. when people log in).
  • The data delivered may be in the same order of magnitude. Lets face it: a web site consists of min 100 items and active readers fetch 0,1 pages per second. This makes 10k x 0,1 x 100 = 100k items/sec. I doubt, that MMOGs deliver more than 100 items/sec to clients. This is 10 times more than what the web server does, but most MMOG items are just IDs of data pre-installed at the client. The web-items are all completely transferred. This reduces the difference. SecondLife transfers really many items by data. For a virtual world it is still an exception and it shows in the bad performance when discovering new areas.
  • Then, there is a stupid reason, why web servers may feel the same load as MMOG servers: average web server technology is much worse, especially if they use scripting languages. Nobody would make MMOG servers in PHP. They are compiled code, C++, C#, compiled Java, compiled Python. But many "Big" web sites run on (uncompiled) byte code. Better than the script, but still byte code.
  • Also, apart from good and bad technologies, there is good and not so good (=stupid) code. This happens to all types of systems and can make an MMOG more responsive than a web site at the same number of users.
  • Across all types of services, some are sharded, some are unique worlds. EVE-Online is unique, which means everyone can play with everyone. WoW is sharded into 1,000 (?) servers of up to 3,000 (?) concurrent players. Your friend is on a different server? you do not even chat to him. Sharding largely reduces communication and DB load. If shard=host then communication overhead disappears. A DB per shard allows for 100 or 1000 times less DB load per instance. This helps a lot.
  • After all, each category covers an order of magnitude. There is much room for "small" or "large" inside a category.
All in all, the classification into "Big", "Massive", "Mainstream" works.

Distributed systems do not really count here. This is about operating the resources required to support concurrent users. E.g. the Web is not operated by a single entity (though you could count the DNS). The Skype network has a (cool) architecture, that lets Skype, the company manage the network without the need to provide all required resources. Skype is also not an issue here. But most systems are operated by someone with a server farm and this is about what these people do.

So, these lessons are for "Big" systems. Actually primarily for web servers but also some messaging issues and general remarks. Some lessons are obvious, some are obvious in hindsight, some are not so obvious, some may be interesting even if you are one of the few (thousand) architects of big systems worldwide. Read the (far from complete but growing) list of lessons here.

17. August 2007

Lessons for Big Systems

Lessons

Take load from the DB

  • Finally the DB is the bottleneck.
  • There is only one DB (cluster), but there can be hundreds of CPUs (web server) and caches (memcache server).
  • Let the CPUs work. 10 web server CPU cycles are better, than 1 DB CPU cycle.
  • Aim at 0,1 DB operations per web page by average.
  • Make it F5-safe. No DB operations for page reloads. No DB for views.
Avoid SQL
  • Keep all live data in memory.
  • Store only for persistency, not for report generation.
  • Use a quick storage, storing 50.000 items per sec is possible
  • DB != SQL, there are quicker interfaces
  • The index is always in memory. That's what SQL DBs are good for.
  • But there are other indexes as well.
External IDs
  • Do not use DB IDs externally. Map all IDs.
  • Use memcache to map external IDs to internal (often DB) IDs.
  • Use memcache as a huge hashtable.
  • External IDs may be strings. After the mapping continue with numbers internally.
DB search loves numbers
  • Everything you search for must be indexed.
  • Avoid indexes on TEXT, VARCHAR. INSERT with index takes significantly longer for text.
  • You may store text in the DB, but do not search for it.
  • You may spend some CPU to map text IDs to numbers for the DB.
100,000 concurrent
  • Imagine 1% of your users are doing the same thing in an instant.
  • If it affects online users, then each task is x 100,000.
  • If it affects all users then everything is x 1-10 Mio.
  • Anything must be at at least 1000/per sec.
  • Do maintenance all the time. There will never be a time of the day where load is so small, that you can cleanup something. Cleanup permanently.
Memcache every business object
  • No object is constructed from the DB.
  • Everything is buffered by the cache.
  • Code with real interfaces, which can be cache-enabled later.
Code for the speed
  • Code for the cache. It is there. It is essential. No way to pretend it is not just for the "beauty" of the code.
  • Write beautiful cache-aware code.
Memcache frontend data
  • Parsing template costs much CPU.
  • Cache generated HTML fragments.
Do not overload the cache
  • Not more than 10 memcache requests per script.
  • If you expect many items, say a mailbos with many messages, then put a summary into a list (mailbox) object even though the same information is in the individual messages.
No statistics on the live system
  • Occasionally they want statistics. Don't do it live.
  • Take snapshots, take the backup. Process it somewhere else.
  • Make statistics offline.
Simple SELECTs
  • Use only simple SELECTs on indexed columns
  • Forbidden keywords: JOIN, ORDER BY
  • Structure and code must guarantee small DB results.
  • Sort in the code not in the DB.
  • If you really need aggregated data, then aggregate permanently. Do not aggregate on demand.
Basics and Trivialities:

Distribute everything
  • Do not rely on a single server for a task.
Check all input

  • Check ALL input.
  • Not only query params are input.
  • Cookies, HTTP header fields are also input.
SQL injection
  • SQL-escape all data in SQL strings.
  • Use prepared statements and variables.
Framework
  • Use a real programming language.
  • Use a compiled language, because the compiler eliminates errors.
  • You will have errors which will wake you at night. So, reduce errors by any means, even if you like script languages.
  • Simple deployment of script languages won't work anyway in the long run, because you will switch on caching and you will have to invalidate the script cache for deployment.