8 ms·
Call me maybe: RabbitMQ
- rdtsc 12y agoIf you didn't read the article, the first part is related to this blog post: http://www.rabbitmq.com/blog/2014/02/19/distributed-semaphores-with-rabbitmq/ http://www.rabbitmq.com/blog/2014/02/19/distributed-semaphor... That post talks about how to build a distributed semaphore using RabbitMQ in the "clustered" configuration option. Aphyr's post is good but RabbitMQ make it pretty clear that this setup is not resilient to network partitions. They even mention it at the end of the blog post. That is also clear in the documentation. https://www.rabbitmq.com/distributed.html https://www.rabbitmq.com/distributed.html The big general question is how likely are you to encounter network partitions? Aphyr believes they are very dangerous and likely to happen to you at some point and should be taken more seriously by the programming world. I tend to take a more of "it depends" on your environment. Recently with VM and cloud deployment network partitions are quite likely. So they should be at the top of your "things to worry about" list. But, I haven't seen it happen on a local LAN. During testing I have induced it by hand (but pulling the network cable out from a switch), but I just haven't seen it happen otherwise. I probably got lucky. Now I have seen other things come like memory corruption, memory leaks in software, and so on. So many other things that worry me more than network partition on a LAN. All that said, it is great to see these experiments run. Please read and study all the other "Call me maybe" series. They are just very good. It turns out most products with a "distributed" component will fall over in the face of a partition. And keep in mind that it is best to have Aphyr discover these issues than your customers ;-)
- hendzen 12y agoAphyr (Kyle Kingsbury) & Peter Bailis have a nice survey of actual network partitions that have occurred in production systems: http://aphyr.com/posts/288-the-network-is-reliable http://aphyr.com/posts/288-the-network-is-reliable
- jamesaguilar 12y agoIf I'm reading the dates right, they admit its lack of resilience to partitions after the Aphyr blog post came out. That's a completely different thing than proactively raising this issue. It's also not clear what if any usefulness a lock service has, beyond as a demonstration, if it can't survive partitions. Network partitions are rare on small scale networks over short durations. But as the saying goes, if you run a billion trials, the one-in-a-million event happens a thousand times. The larger a network grows, the more likely you'll eventually experience the joy of a partition.
- rdtsc 12y ago> they admit its lack of resilience to partitions after the Aphyr blog post came out. That's a completely different thing than proactively raising this issue. In that blog post. This doc page was there long before. https://www.rabbitmq.com/distributed.html https://www.rabbitmq.com/distributed.html Network partition tolerance is a general configuration trade-off not related just to that particular trick of distributed semaphores. Presumably someone who set up a clustered configuration already decided on the likely-hood of experiencing a network partition and read the docs. Now there is a another issue explored and that is tolerance to network failures in general between clients and even a single server. That is (the way I understand it) not related to clustered or un-clustered configuration. It relates to the stability of network connections between clients and server(s).
- jamesaguilar 12y agoI don't know if you read the blog post, but it demonstrates that the partition tolerant mode from the page you just linked is not actually, and that RabbitMQ, even in that mode, can't be used as a lock service. I didn't see the part of the post about connections between clients and servers, but I was reading on mobile, so maybe I just missed it.
- rdtsc 12y ago> but it demonstrates that the partition tolerant mode from the page you just linked is not actually I was talking about the part about picking a response to a partition detection event. In this case "ignore" (see "Recommendations" section), instead of "auto-heal" or "minority-pause".
- Nitramp 12y agoI think this is a misunderstanding in terminology. "Network partition" in this context does not necessarily mean your LAN cable is broken, it means that you have a cluster of servers that expect to communicate, and some of them cannot communicate between each other. That could be because of a faulty network, but it could also be due to hardware failures, misconfiguration, slowness of servers (hard to tell from not reachable), and a whole host of other causes. If you have more than just a couple of servers, this will happen sooner or later. If you only have a couple of servers, IMHO don't bother with a truly distributed system, use a database with failover to a replication slave (lazy or eager, depending on what consistency you need) and be done with it.
- ams6110 12y agoI've seen it happen, not what I'd call often but not so rare that you're astonished by it. NICs fail, switch ports or entire switches fail, stuff like that happens from time to time.
- Tomdarkness 12y agoIn relation to the second part of the article regarding lost messages it would be interesting to see the same test done but with HA queues: http://www.rabbitmq.com/ha.html http://www.rabbitmq.com/ha.html I might be wrong but I'd suspect that it might solve the problem as while the partioned nodes will still discard the data it should be present on the master node.
- jpgvm 12y agoThose tests were done with mirrored queues: "We’ll be using durable, triple-mirrored writes across our five-node cluster"
- KayEss 12y agoThe seriousness of this depends on the consequences of the failure and how much of the time the lock is held. If the failure causes double charging a customer that's likely far more serious than if it double emails a customers. And if the lock spends most of its time not held then spotting a doubling of its message can probably be spotted and fixed far more easily than if the lock spends most of its time held.
- lobster_johnson 12y agoRabbitMQ requiring a reliable network is causing problems for us in production. Anyone else struggled with this? We're running several clusters on different providers; one is Digital Ocean, another is on a partner's VMware vMotion-based system. The kernel gets soft lockups now and then (in the case of vMotion, when VMs are automatically migrated to other physical nodes), which causes RabbitMQ partitions. The lockups may last a few seconds, but I've seen minute-long pauses. When this happens, RabbitMQ starts throwing errors at clients. Queues disappear and clients can't do anything even though the local node is running. Although I understand the behaviour, that's not what I want from a queue; I want a local node to store messages until it rejoins the cluster and can dispatch them, and I want the local node to continue offering messages to those listening. Unfortunately, RabbitMQ's design doesn't allow this: Queues live on the node they were created on, and RabbitMQ does not replicate them. We have turned on autohealing, but I don't like the fact that minority nodes simply wipe their data when they rejoin the cluster. The federation plugin doesn't look like a great solution. I really like RabbitMQ, but maybe it's time to considering something else. Any suggestions? Something equally lightweight and "multimaster", but without the partition problem?
- jpgvm 12y agoYou need to look into using HA queues (ie. mirroring) and ensure you are using a client that properly supports consumer cancel notifications and are actually processing them correctly. There is also significant TCP/IP and RabbitMQ tuning that can be done to make fail-over much faster etc. I have used RabbitMQ for years and though it's not perfect if you do try something else you will quickly understand that it's still so far ahead of the competition. Best of luck!
- atombender 12y agoHA queues have this caveat in the documentation: This solution requires a RabbitMQ cluster, which means that it will not cope seamlessly with network partitions within the cluster and, for that reason, is not recommended for use across a WAN And then: However, there is currently no way for a slave to know whether or not its queue contents have diverged from the master to which it is rejoining (this could happen during a network partition, for example). As such, when a slave rejoins a mirrored queue, it throws away any durable local contents it already has and starts empty. So I don't think that's helpful to us at all. We are going to migrate to a client that supports consumer cancel notifications, though. Thanks for the tip.
- nateburke 12y agoCool. I've been bitten by RabbitMQ before. When are you going to give cassandra the jepsen treatment?
- jpgvm 12y agoThe unfortunate reality for RabbitMQ is about 99.99% of the problems people have with it (assuming they are using it as a queue) is that they either a) are using a bad/incomplete client or b) aren't using the HA features correctly or c) are doing all of that but then aren't checking for failure states (CCN etc) properly so their application isn't responsive enough in failure conditions.
- slantedview 12y agoI don't think any of that applies here.
- shadytrees 12y agoIt already exists! Back from the 2013 series. http://aphyr.com/posts/294-call-me-maybe-cassandra/ http://aphyr.com/posts/294-call-me-maybe-cassandra/
- mjamil 12y agoA bit meta. That was a really well-explained description of a non-trivial problem. I wish there was (or I knew of) a curated collection of bloggers that wrote this well about software.
- steveklabnik 12y agoThe rest of Aphyr's stuff is just as good. My personal favourite: http://aphyr.com/posts/283-call-me-maybe-redis http://aphyr.com/posts/283-call-me-maybe-redis
- freshhawk 12y agoAphyr's Clojure tutorial series is one of the best programming language tutorials I've ever read. I'm definitely a fan. http://aphyr.com/tags/Clojure-from-the-ground-up http://aphyr.com/tags/Clojure-from-the-ground-up
- EGreg 12y agoHow does RabbitMQ compare with 0MQ?
- lobster_johnson 12y agoThey are conceptually very different. Rabbit is a broker based on the AMQP model of exchanges, queues and bindings. Producers and consumers talk to Rabbit, which processed, stores and dispatches messages across the cluster. It supports both direct (one message goes to a single consumer) and fan-out (every message goes to all consumers) delivery, and can pick targets based on wildcard routes. Queues can be "durable", where Rabbit will store messages on disk atomically in case of a crash. You can restart Rabbit, as well as producers and consumers, and the queues will still be intact. ZeroMQ is (despite the name) not a message queue in the classical sense; it's more like multicast TCP or UDP, with some message queue functionality. To start, since Zero is embedded in each program, a "queue" does not persist if a producer and consumer are not both running; any buffering occurs in the consumer, so if a producer emits messages and no consumers are running, the messages are lost. Similarly, there is no durability; all messages must be handled right away. ZeroMQ is a low-level toolset for implementing robust messaging -- it could be used to implement a high-level broker like Rabbit, but in itself it doesn't compete directly.
- keypusher 12y agoI thought that in order to create a distributed locking system, you need to be able to reliably fence unreachable nodes. "Network Partitions" sound a lot like "Split Brain" to me. I am more familiar with traditional clustering solutions such as Pacemaker/Corosync and the GFS2 DLM than RabbitMQ, so perhaps I am missing something here?
- teraflop 12y agoThat's more of a software design consequence than a fundamental limitation. It's common to build fault-tolerant systems by designing a master-slave architecture, and then bolting a failover mechanism on top so that exactly one node is the master at any time. That approach suffers from exactly the problem you describe: if you can't contact a node, there's no foolproof way to ensure it isn't acting as a master, so you risk split-brain unless you can remotely fence it. Consensus algorithms like Paxos/Raft don't suffer from this problem; in order to make steady progress, there should be only one master, but they still operate consistently if that assumption is violated in the short term. So no fencing is needed. A lot of people seem to have a fundamental misunderstanding that using Paxos for leader election is enough to make a system consistent, and it's emphatically not. For instance, I hope this isn't still true in current versions, but for a long time HBase was vulnerable to losing committed data because it misused Zookeeper and allowed multiple regionservers to simultaneously "own" a given table region. (I found about that issue when it happened to me in production, so now my default assumption is that all "distributed locking" systems are broken until proven otherwise.)
- deleted 12y ago[deleted]
- deleted 12y ago[deleted]
- deleted 12y ago[deleted]