3 ms·
I'm not sure I understand the issue here: If you have an at-least-once queue, it might still be reasonable for the application to directly write to a transactio
by bcoates 12y ago
I'm not sure I understand the issue here: If you have an at-least-once queue, it might still be reasonable for the application to directly write to a transaction log somewhere that can be reconciled for auditing, but that audit should always pass (even under fault conditions): there is either a design flaw in the system or malicious insider tampering if the queue does not deliver the message. Is it necessary to audit for that sort of thing on a sub-day level?
It doesn't seem like an appropriate place to put a wait-and-retry loop -- it absolutely breaks the entire point of queueing if the sender does not consider the task completed and flush all of its state once the message has been successfully enqueued.
- ryanjshaw 12y ago> there is either a design flaw in the system or malicious insider tampering if the queue does not deliver the message Aside from the partitioning/client issues that hueyp points to, queue managers with transient storage can crash or lose power. Queues managers with persistent storage can suffer corruption or data loss. In this way the 'at least once' "guarantee" will not be met, on the queue side. I have noticed a model forming these days where you have lots and lots of queue managers that mirror each other and provide a sort of defense-in-depth against failure. It's an interesting approach, and very convenient, but infeasible on a low budget. I think you also need to be very careful with such an architecture because it seems like it could be easy to accidentally break it.
- bcoates 12y agoAt-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?
- hueyp 12y agohttp://queue.acm.org/detail.cfm?id=2187821 http://queue.acm.org/detail.cfm?id=2187821 Section: "Zero or more times ... gauranteed" "When considering the behavior of the underlying message transport, it is best to remember what is promised. Each message is guaranteed to be delivered zero or more times! That is a guarantee you can count on. There is a lovely probability spike showing that most messages are delivered one time. If you don’t assume that the underlying transport may drop or repeat messages, then you will have latent bugs in your application. More interesting is the question of how much help the plumbing layered on top of the transport can give you. If the communicating applications run on top of plumbing that shares common abstractions for messaging, some help may exist. In most environments, the app must cope with this issue by itself." At least once, at most once, etc are just probability distributions. All you can count on is zero or more times.