13 ms·
There Is No Now – Problems with Simultaneity in Distributed Systems
- nickbauman 12y agoThis article underscores a point I always end up explaining to people who look at Cloud Computing (implemented as "the lambda architecture" in best-of breed scenarios) as a golden hammer that it's a good technology, but only for a certain class of problems. You can do things like monitor trends over time and even act on them with soft deadlines using Cloud. But you will never have a cloud technology control your anti-lock braking system on your car. It basic understanding of CAP theorem, really.
- dragonwriter 12y agoWell you could have locally accessed dynamically provisioned computing capacity (cloud technology) in your car doing that. Remote/shared hosting is a different and older thing than cloud technology (though it's a popular use of cloud technology.)
- jerf 12y agoYou already do. "Multitasking" covered that decades ago. I'm not sure that's a useful metaphor to cloud technology, though, most especially because multiprocessing systems are exactly where we got our wrong impressions about reliability and the lack of network issues in the first place....
- Retric 12y agoI would be shocked if any ABS systems in a car actually dynamicly shared resources with anything else.
- padelt 12y agoThat is not uncommon at all and not unsafe if done right. Hard realtime systems are there for exactly that reason. Since you need to dimension everything to worst case scenarios, gains tend to not be as big as with soft realtime systems but there you have it: Dynamically shared resources in critical systems.
- Retric 12y agoI can see Vehicle Stability Assistance systems that also do ABS breaking, but that's not really dynamic resource allocation so much as new solution. Ditto, for sending information over the cars local network. Can you give an example where something not related to traction is run on the same CPU as the ABS/traction system?
- pjc50 12y agolocally accessed dynamically provisioned This is a mission critical safety system. There are some standards which are unhappy with dynamically allocating memory in that situation, let alone dynamically allocating the entire compute resource.
- dragonwriter 12y agoThat's a good point, though it is different consideration from the response time issue imposed by physical limits of communications round-trip time.
- jgrahamc 12y agoRear Admiral Grace Hopper (one of the most important pioneers in our field, whose achievements include creating the first compiler) used to illustrate this point by giving each of her students a piece of wire 11.8 inches long, the maximum distance that electricity can travel in one nanosecond. I made download and print versions of this: http://blog.jgc.org/2012/10/a-downloadable-nanosecond.html http://blog.jgc.org/2012/10/a-downloadable-nanosecond.html
- glittershark 12y agoReminds me of this (absolutely spectacular) talk from 2009 by Rich Hickey: http://www.infoq.com/presentations/Are-We-There-Yet-Rich-Hickey http://www.infoq.com/presentations/Are-We-There-Yet-Rich-Hic...
- rdtsc 12y agoKeeping an always consistent state in a large distributed you are fighting against the law of physics. Like it is mentioned, Google did it with their F1/Spanner SQL database. But that also mean GPS receivers with antennas on the roofs of data centers. Which is yet another thing that can fail and it, either by itself or in a cascade of other failures will lead to unspecified and possibly undesirable behavior. Recently I see a lot of advocates of dropping NoSQL databases and moving back to Postgress or other SQL databases. The problem is SQL and schemas is not the only reason NoSQL databases became popular, they became popular because they also started to have a default and more defined behavior in respect to replication and distribution. Most solutions don't need that and sticking with a solid single database works very well. But those that need distributed operation have a pretty hard task ahead of them. One heuristic you can look at is if and how is distribution implemented. Is it something that is bolted on top, like you download some proxy or addon, or added as an afterthought, be very careful. Those things should be baked into the core. For example, does it support CRTDs? Does it have master to master replication and so on. If it claims to have consistency instead availability figure out how it is implemented, Paxos, Raft, or something else. So far I think Riak is probably the database that thought the hardest and did the most reasonable job. The other simpler database is CouchDB, they have well specified conflict resolution behavior and master to master replication. But you have to usually write your own cluster topology. There are probably others but those are the two I know of first hand.
- manigandham 12y agoNice points, have to mention Aerospike as another that did a lot of the right stuff with regard to distributed actions and clustering.
- logicallee 12y agoI'm not an expert, and it sounds like you are, so I appreciate your feedback here: what do you even mean by a consistent state? even in theory a person initiating a new additional record in Auckland, New Zealand at the same time somebody iniatiates a change in Gibraltar or London (which are antipodal to the former[1]) 66 milliseconds away, cannot have a confirmation in less than 120 milliseconds, right? So do you just wait for that before declaring 'consistency'? Do you literally add 120 milliseconds to each and every request? (And this is assuming you have a damned good solution to the two generals problem)? I mean suppose the database tracks something as simple as: number of web page hits. It's a counter. You now distribute it, and have a stochastic process of counter hits between 0 and 5 per second in your largest cities, distributed throughout the world. How can that database ever be consistent? If there are ten new records per second in New Zealand and ten records per second in the UK, and they potentially depend on each other in some way, are you going to just make everyone wait until everything has been committed and confirmed to be consistent? Or is "a foolish consistency the hobgoblin of little minds", and you really can accept out-of-date data and deal with merging conflicts later? I just don't understand why we would expect consistency to rank up there, when we deal with a worldwide real-time system where the difference between getting served by a local database in 40 milliesconds and one far away in 250 milliseconds is both staggering, and incredibly noticeable. why be consistent? what is consistence? [1] http://www.findlatitudeandlongitude.com/antipode-map/ http://www.findlatitudeandlongitude.com/antipode-map/
- tlarkworthy 12y agoSpanner is externally serializable so you do get a 'now'. You just don't know what the agreed now was until after the write. The idea that there is no 'now' is of course, preposterous, we have very strict laws of physics supporting the concept of now (sans relativity) and eventually our engineering will be able to track that very accurately. Spanner is a step on that journey. These kinds of 'impossible' articles will appear very dated in 10 years time, as they are really over exaggerating the rules-of-thumb of the previous 10 years.
- dragonwriter 12y ago> we have very strict laws of physics supporting the concept of now (sans relativity) Laws of physics "sans relativity" aren't the actual laws of physics in our universe. There very much is no now except the now that is also here. Its quite accurate to say that simultaneity does not exist in distributed systems, and simultaneity is less valid as even an approximation the more widely distributed a system is.
- nostrademons 12y agoTo add some context & numbers to dragonwriter's point, the speed-of-light delay from New York to San Francisco is about 21ms [1]. This is about 5 disk seeks, 1200 random SSD reads, 200K main memory reads (without caching), or 10M CPU cycles [2]. Speed of light delays absolutely matter in a distributed system. [1] http://chimera.labs.oreilly.com/books/1230000000545/ch01.html#PROPAGATION_LATENCY http://chimera.labs.oreilly.com/books/1230000000545/ch01.htm... [2] http://www.eecs.berkeley.edu/~rcs/research/interactive_latency.html http://www.eecs.berkeley.edu/~rcs/research/interactive_laten...
- tlarkworthy 12y agothat's transmission delays not relativistic effects. The idea of spanner is you have a timestamp of when it was decided to commit. The quorum knows they can't contact each other quickly but they trust each others timestamps and resolve conflicts based on the trusted commit times (which are accurate through hardware). Transmission delays don't undermine the fact there is a very real concept of ordered time in the physical world which is exploitable[1]. Spanner exploits it faster than transmission delays, but with a clock error of 10 ms or something. We have better clocks than that so its probably going to improve... [1] sans relativity effects which are TINY, and not the limiting factor at the moment.
- jamespitts 12y agoThis is an excellent overview. Simultaneity isn't just for machines -- it is necessary for people being connected together online as well. Time sync is a huge part of creating a shared experience, and this will become more widely appreciated as virtual reality develops socially. We solved this in a very limited way at rapt.fm for our timed rap battles. We maintained a shared clock tick (with adjustment), allowing UI and in-game events -- e.g. the beat kicking in -- to happen somewhat simultaneously across browsers. This helped make up for the latency of video, and created a feeling that people were together at the same "place".
- brianzelip 12y ago:) for the dude eating from a bowl of cereal at the end of the rapt.fm landing page background vid
- jamespitts 12y agoYes, that is me! We had a lot of good times in downtown Det.
- jhayward 12y agoThanks, this is a helpful introduction to the history and literature of concepts related to time in distributed systems. Most people's concept of time is quite simple and they need to be broken loose of some intuitively held, but unhelpful beliefs before they can really do engineering with respect to time. Is there a paper somewhere that new folk should read first? One that includes: - A tutorial to describe all of the things we think of as 'time', e.g. order of sequence, etc. and their dependence on each other. - The idea that time as it occurs in the physical world is probabilistic - requiring such descriptors as precision (what is the smallest difference we can discern), and error bounds or probability distributions (how accurately can we describe it). - And for the concrete thinkers who 'get' that true simultaneity is impossible, an easy to understand example of how we succeed in observing logical coherence, from the scale of a single CPU chip (internally non-coherent, externally consistent) to cross-continent compute clusters?
- nuxi 12y agoThere was a great presentation given by Tom Van Baak on measuring time, precision, accuracy etc. at this year's FOSDEM . Abstract an video are available at https://fosdem.org/2015/schedule/event/precise_time/ https://fosdem.org/2015/schedule/event/precise_time/, he has lots of additional time-related information on his website http://leapsecond.com/ http://leapsecond.com/.
- stcredzero 12y agoWith the way the world looks now, and the way it's shaping up to be in the future, we should be thinking about designing systems where there is no "now" but rather, there is instead a notion of "everything syncs soon." There can be a robust consensus reality in such a system, but it has to exist a short interval of time in the past. I'm currently working on a multiplayer game design on these principles.
- jchrisa 12y agoIf you want to let a mature open source sync system do the heavy lifting we've recently added Unity 3D support to Couchbase Lite. http://developer.couchbase.com/mobile/unity/ http://developer.couchbase.com/mobile/unity/ Hello world is a game where players drop shapes into a world and they automatically sync / appear on other devices. Sorry about the plug -- great article Justin! I sent it to my team because I think it does a great job stepping back from the buzzwords and actually talking about the problems distributed databases have to confront.
- stcredzero 12y agoInteresting. However I'm not just looking for syncing stats and inventory. What I'm talking about is designing a new game mechanic around the realities of distributed gaming, then designing a more robust server around that. The server will be finalizing entity positions, damage, and collisions a half second behind "real-time," and the game mechanic will be designed around this in a way that preserves player autonomy and immediate response. (Basically, deemphasizing aiming while emphasizing tactical dodging where the player seeks to minimize received damage over time.) great job stepping back from the buzzwords Those pesky buzzwords. You were just caught out being misled by one. (sync -- which typically means what you thought it did, but also has a generalized meaning)
- marknadal 12y ago@stcredzero I'm interested in building simulation systems for distributed systems, and trying to "gamify" them to teach people about these concepts in a tactile-playable way. Would you be interested in such work? It is less building a "entertainment" game and more an "educational" fun game. Shoot me an email: mark@gunDB.io
- marknadal 12y agoHallelujah, we need more articles like this. I work in distributed systems, and people hate hate hate hearing the truth, let's summarize the article: 1. You cannot beat the speed of light. 2. Machines break. Even the most reliable ones. 3. Networks are unreliable. Even local area networks. 4. It is an exciting time for distributed systems: CRDTs, Hybrid Logical Clocks, Zookeeper, etc. These are things I feel like I've been preaching a lot, but get upset responses "Well, certainly Globally Consistent systems work, most databases are!" Lies, as @rdtsc mentions: "Keeping an always consistent state in a large distributed [system] you are fighting against the laws of physics." Next up, that is why master-master replication is important, because at some point your primary will go down (or at least the network to it). I started an interesting thought experiment that turned into a full on open source project: What if we were to build a database in the worst possible environment (the browser, aka javascript and unreliability)? What algorithms would we need to use to make such a system still survive and work? This is why I chose and worked on solutions involving Hybrid Logical Clocks and CRDTs, which are at the core of my http://github.com/amark/gun http://github.com/amark/gun database. An AP system with eventual consistency and no notion of "now", as every replica runs in its own special-relativity "state machine" view of the world. These are all interesting concepts, and the article was a good one. I recommend it.
- vdm 12y agoNice to see another example of applying CRDTs in the browser, I only knew about swarm.js. http://swarmjs.github.io/ http://swarmjs.github.io/ They have a few good posts on why you would want to do this.
- keppy 12y agoI suggest reading this article for some insight if you haven't built out a distributed system. Good hands on, practical information: http://engineering.linkedin.com/distributed-systems/log-what-every-software-engineer-should-know-about-real-time-datas-unifying http://engineering.linkedin.com/distributed-systems/log-what...
- aristus 12y agoRelated: http://carlos.bueno.org/2010/04/dismal-guide-to-concurrency.html http://carlos.bueno.org/2010/04/dismal-guide-to-concurrency....
- ajarmst 12y agoDoesn't appear to mention Lamport Timestamps (https://en.wikipedia.org/wiki/Lamport_timestamps https://en.wikipedia.org/wiki/Lamport_timestamps) (Yes, THAT Lamport), which are one of the most elegant mechanisms to deal with the some of the discussed problems.
- macintux 12y ago> Another such area of work is logical time, manifest as vector clocks, version vectors, and other ways of abstracting over the ordering of events. Vector clocks and version vectors are variations on the Lamport clock concept (and that quote is from a paragraph that mentions Paxos, another Lamport invention, cited in the bibliography).
- keppy 12y ago> Another such area of work is logical time, manifest as vector clocks, version vectors, and other ways of abstracting over the ordering of events. This idea generally acknowledges the inability to assume synchronized clocks and builds notions of ordering for a world in which clocks are entirely unreliable. Hardware is unreliable. Software is possibly less reliable. We have known that for a long time. The author talks on a conceptual level about logical time, but this concept isn't enough to understand the real challenges & possible solutions of keeping interactions in your system logically ordered in the dimension of time[0]. > You can think of coordination as providing a logical surrogate for "now." When used in that way, however, these protocols have a cost, resulting from something they all fundamentally have in common: constant communication. For example, if you coordinate an ordering for all of the things that happen in your distributed system, then at best you are able to provide a response latency no less than the round-trip time (two sequential message deliveries) inside that system. Consensus protocols don't provide a logical surrogate for 'now', a log does that. The silver bullet for assuring that your transactions are ordered correctly is immutability[1]--"If two identical, deterministic processes begin in the same state and get the same inputs in the same order, they will produce the same output and end in the same state.[0]" It's important, from the perspective of the implementor, to understand that there are multiple pieces to this puzzle, and that each protocol has very specific details that can make or break the reliability and performance of a distributed system. This is similar to how a small bug in your cryptography code can expose the entire system to threat. Paxos itself can be implemented in a myriad of ways, and each decision the implementor makes must be well researched. [0]http://engineering.linkedin.com/distributed-systems/log-what-every-software-engineer-should-know-about-real-time-datas-unifying http://engineering.linkedin.com/distributed-systems/log-what... [1]http://basho.com/clocks-are-bad-or-welcome-to-distributed-systems/ http://basho.com/clocks-are-bad-or-welcome-to-distributed-sy...
- yodsanklai 12y ago> One of the most important results in the theory of distributed systems is an impossibility result, showing one of the limits of the ability to build systems that work in a world where things can fail. Layman question. I wonder how important this result really is. This is an impossibility result in a certain model, where processes are deterministic. It's certainly a nice theoretical result but in practice, there are probabilistic algorithm that solve this problem. I don't know what are the probabilistic bounds of probabilistic consensus algorithms, but if it's arbitrary low, the impossibility result for deterministic processes is irrelevant isn't it? After all, if we can live with a super low probability of a meteorite destroying the planet, so can we with a good probabilistic consensus algorithm.
- hyperion2010 12y agoI really enjoy the perspective that "now" is often a useful abstraction for certain types of processes. The fact that it turns out that "now" is one hellishly leaky abstraction. My perspective coming from biology is that for many systems the only meaningful type of "clock" is a logical clock. The important thing is not "when," when is used as a proxy for an assumed state of a remote part of the system (even if remote is only 10cm away), logical clocks are the only source that can guarantee that the state of the system is what you expect it to be so that it will perform as expected. Thanks to the many hardware guys who have spent years working out the underlying logic for this we mostly ignore it for things like processors. Now we just need to solve it for arbitrarily large finite delays! This also reminds me of a very funny (or depressing) read on systems engineering by James Mickens [1]. 1. http://research.microsoft.com/en-us/people/mickens/thenightwatch.pdf http://research.microsoft.com/en-us/people/mickens/thenightw...
- Terr_ 12y agoAnother fun exercise is dealing with players who say "OMG teh netcode sux" for online first-person shooters, especially when they have patently unrealistic expectations for how well the software should break the speed of light. Sometimes the hardest part is getting them to understand exactly how much of what they take for granted is an illusion... Often even before any packets leave their machine.
- Dylan16807 12y agoIt's less speed of light at fault than other factors, though. Bufferbloat is terrible, and it's hard to make netcode that reacts well when a few percent of packets are lost or slow.
- Terr_ 12y agoSure, other factors predominate, but even in a significantly-more-ideal world, a signal from LA to NY is still going to be a 40ms round-trip. (3940km one-way along the surface of the earth, 200,000km/s signal speed in the glass.) That's still more than enough time to require algorithmic trickery from games in order to provide the illusion of "real-time" gaming over the internet.
- Dylan16807 12y agoThink about how many console games run at 30fps. Then consider that a fully framebuffered game running at 30fps, even under ideal conditions, has 70ms of latency by the time a frame finishes displaying. 40ms isn't a big deal.
- Terr_ 12y agoEven assuming a perfect computer with zero input latency and zero display latency the network part still matters when it comes to knitting an area like "North America" into a game region. It feeds into longer causal chains. For example, consider two players in LA both using an NY server over the dedicated zero-overhead fiber-optic network mentioned above, with a relatively modern game-networking stack. (1) Player A raises an energy-shield t=0 ms. (Client-prediction.) (2) The server agrees as of t=20. (3) Player B shot an instant-travel hitscan weapon at t=39, while the victim was still exposed on his screen. (4) The server gets the shot-message at t=59, and honors it (Latency Compensation), sending a damage/death message out to Player A. (5) Player A receives the news of his damage/death at t=79. Even in that world of unattainably-good equipment, that's 80ms of "wait, that doesn't look right".
- jbergens 12y agoOne piece of warning from the article regarding "last write wins": >..., this is really a "many writes, chosen unpredictably, will be lost" policy—but that wouldn't sell as many databases, would it?