3 ms·
At-least-once[1] queues are sound under bounded partitioning and can withstand any specified amount of data loss[2]. Uncorrected data corruption errors cause wh
by bcoates 12y ago
At-least-once[1] queues are sound under bounded partitioning and can withstand any specified amount of data loss[2]. Uncorrected data corruption errors cause whole-system failure in non-distributed systems as well as distributed ones and the solutions are the same in both.
You're not just throwing a bunch of queues at the problem and hoping they stick. Your queue based system is in one of two states: it has a specified amount of fault tolerance based on its design, or it is defective (see the 'call me maybe' posts.) If you have a defective distributed system adding more nodes makes it worse, so it is vitally important to know which one you have and design accordingly.
[1] Calling them zero or more times queues is profoundly bullshit in that it conflates faults (multiple delivery) and failures (non-delivery), which is reliability 101. Failures are expected to propagate; faults must not, and a design which attempts to work around failure instead of just faults is almost certainly a sign of deep confusion.
[2] This doesn't violate CAP because a client must write to enough nodes that it can be certain that it is not part of a doomed partition (potential loss of availability), or pretend the write succeeded without that guarantee (potential loss of consistency).
- ryanjshaw 12y agoIf I've understood you correctly, you're saying that when designing a distributed system, (1) non-defective database(s) and queue(s) have some degree of fault tolerance, which we can assume for argument's sake is the same (whatever we can do to make our database reliable we can do to make our queue reliable), and (2) whichever fails - database or queue - the resolution is the same in both? If this understanding is correct, the point of contention is #2. Scenario | Untracked | Tracked | Queue data loss | (A) | (B) | Database loss | (C) | (C) | (A) we can't replay events from database to queue (B) we can replay events from database; (C) we have to restore from backup Obviously the two scenarios are indeed the same if we lose the database (C). But when it comes to losing the queue, if we're tracking what we've put on the queue and whether we've received an acknowledgement or not, we can easily replay events to that queue (B). If we didn't track what we put on the queue, how do you recover in (A)? As I see it, the database represents a consistent view of the application state and data -- e.g. "payment 123 was settled". If you lose your database, you have to restore, then move upstream and replay events there (if possible). Downstream systems handle duplicates in replay scenarios by discarding them. This is a nice straightforward way to recover. If you lose your queue, it's not an issue, because the queue is just a channel and you replay events over it. What's important is your higher-level protocol across the channel. In my mind the pattern is: S -> channel -> R. Channel can be a brokered queue, a brokerless standalone queue, async IPC, a thread delivering messages ala ZMQ, a file, a web API, raw UDP, whatever, but either way its the same pattern -- S tracks its sending state, it doesn't trust the channel because if things break you don't have a consistent view of your downstream obligations that you can recover from. The channel's job is to deal with delivery, but it can't address higher-level concerns across the channel. We put messages onto queues as part of a higher level abstraction -- we use TCP/IP, for instance, as part of an HTTP conversation: if a web browser sends a GET and doesn't get a response, it doesn't just sit waiting forever hoping the TCP/IP stack will figure things out... It recovers and maybe even retries but warns you about the problem so that you notice you typed the wrong IP address in and fix it. TCP/IP is a nice unreliable example, but whatever your channel is, it's going to be unreliable, so in our business applications we need a business-case specific protocol to record our business transactional state as well to handle the failures. If you offload that obligation to the channel, you're stuffed because the channel can't understand your business case and consequently how will you recover? Unless you mean you can design systems with a degree of fault tolerance so great that its not worth considering failures? That can't be right, what am I missing?