You cannot have exactly-once delivery (2015)
bravenewgeek.com
bravenewgeek.com
The thing is, having a print company that gets the printing part right is more important than having one that gets the API right. I use them anyway, and accept the risk that there will very occasionally be duplicate orders. At least in my business, it's just tshirts.
A few months ago I bought a fairly expensive cordless vacuum from hoover.com. I was charged once, but two of them arrived. I suspect I know why.
What did they suggest you should do if you sent the order once and it didn't arrive?
It's something I have been tracking over the years, rare but happens. Again some non transactional, non idempotent integrations setup out there...
Oh, there's a fourth kind: "none-of-the-above", i.e., neither at-most-once or at-least-once. The message gets delivered between [0, ∞] times. Maybe it gets delivered … maybe not. Your message is like UDP packet.
A surprising number of systems exhibit this behavior, sadly.
I noticed [0, ∞] delivery semantics in a widely-used, internal/homegrown message delivery system at a big tech company once. The bug was easy to spot in the source code (which I was looking at for unrelated reasons), but the catch-22 is that engineers with the skills to notice these sorts of subtle but significant infra bugs are the same engineers who would've advised against building (or continuing to use) your own message delivery system in the first place when there are perfectly serviceable open source and SaaS options.
I don’t think it’s as simple as always preferring to bring in external things.
So, if you look at it juuuust right, ELK can be considered as delivering messages … log messages. Normally our ELK system exhibited the behavior I described above: our logging component would make best efforts to ensure that log messages did get delivered. But there was no ID on the log message when it was submitted: double-submission would result in the message getting duplicated. The message was only removed from the queue of messages that needed to be submitted if the ACK was successfully received. The local queue was only so big: if the application continued to log but couldn't submit to ELK, well, it would just discard the messages², so messages might not get delivered period, even when the client was fine, (e.g., during network outages).¹
That was all fine and good (the consequences of a log message getting delivered twice or not delivered under normal circumstances is "whatever").
One day, the team managing ELK misconfigured it, causing all log submission to start failing. This … didn't get detected? (I mean, as a consumer of the service, what are you going to do, log the error?) Worse, some logs were first routed through the local syslog daemon. It decided to log that it couldn't send the message to a local log file, and then retry without backup, and never gave up. (And, of course, there's not log rotation on that local log file.) So we noticed the problem when the disk filled up … and then realized basically every VM in the fleet was doing this. So it felt like ∞] that day.
But yes, mathematically, that bound should be ∞).
¹there is some decent discussion in the thread below my OP about whether this technically counts as "0 deliveries".
²the logic being that crashing due to inability to log is not worth it. It noted the failure on its stderr … but you had to know to look.
First rule of networking: Every bad thing that can happen, will.
If nothing else, by sheer bugginess you will certainly have something, at some point, retransmitting UDP packets for no good reason.
I think that is the point of these conversations. You shouldn't have application expectations that cannot be met in the real world. Even in your own data center you have far less control over your mirrored switches doing something dumb like sending a stream of packets twice out of their respective interfaces.
> There are essentially three types of delivery semantics: at-most-once, at-least-once, and exactly-once. Of the three, the first two are feasible and widely used.
But in the context you're talking about, even those first two are not "feasible." And in the context you're talking about, there's basically no point in talking about anything you can do to try to achieve certain behavior since there is technically no way to physically guarantee that. What you are saying is definitely not "the point of these conversations," because you're just saying that literally nothing can ever be absolutely guaranteed. Technically true, but not particularly helpful when you're designing real systems to solve real problems, and your systems will mostly operate on infrastructure that mostly does work as intended.
In that context, everything bad can happen, but it rarely does. For example: packets most definitely get dropped, but not every packet does in practice. Packets can get double-sent, but again, not every packet in practice. Packets can get corrupted randomly, but again, not every packet in practice.
The hard mathematical guarantee: nothing at all! You will get a packet from 0 to N times, where N might be arbitrarily large.
The in-practice behavior: mostly kinda works, but be careful! Have a strategy for reordering, a strategy for detecting and handling dupes, and a strategy for when you need to just give up and start from scratch.
The internet is like that but if some router is screwing it up, you can’t call the owner and complain.
Nowadays though, it is getting much harder to do so however. I know there's a secret IRC of black belt NetOps out there, but I haven't managed to figure out how to route myself there yet.
(Rumor has it there's a router out there with a reliably flaky network card that'll mangle the packets juuuuuust right. Personally, I think it's all hooey and someone has a really neat IPTables file complete with port knocking)
It turns out it really isn't that much daylight, and it's really easy to overestimate the size of that window. As you scale up, that window gets smaller and smaller, too, which makes it even more exciting.
Nevertheless, it is indeed where real systems are built.
I also disagree that realizing that nothing can be guaranteed is not important to building systems. The more you scale up, the more important it is. It's hugely important. It's one of the major things that separates people who can build real network systems at scale and those who can't. Far from the only such thing, but certainly one of the important ones.
On public networks, yes.
Of course in many cases andyou have to limit what edge cases you deal with unless you have infinite development time, but unexpectedly repeated UDP packets are definitely something that happens often enough to account for it if your protocol could be adversely affected by it.
Packets can and will be queued in multiple outgoing interfaces.
Dumb shit happens in network kit.
Also, the new SDN stuff sends packets multiple times over different paths on purpose. It's supposed to discard everything except the one that got there first, but...
On noisy wifi you're transferring data, and the destination finally gets enough packets to send an ACK for the sliding window, only the source never gets the ACK, so it sends the packets it thinks you want but already have. Some of those get through, and the destination realizes it needs to send the ACK again because clearly you didn't get it the first time. Finally you resync and start getting new data, until the next cup of coffee goes into the microwave and it all repeats again.
Since many UDP protocols end up re-implementing half of TCP, you're going to have some of the same failure modes.
They're supposed to run STP or bonding or something to make 1 & 2 a logical single link, but that's misconfigured, misbehaving, or just plain old buggy.
By sheer bad luck, switch B has currently overflowed its MAC lookup table, and is falling back to broadcast for your destination MAC.
You send a packet to switch A. Switch A looks up the destination MAC, and forwards the packet to link 1. Switch B receives the packet, has no idea where that MAC is, and forwards it to links 2-n. Switch A receives the packet, looks up the destination MAC, and and forwards the packet to link 1. Rinse, repeat. Observe packets sent by switch B on ports 3-n.
> The way we achieve exactly-once delivery in practice is by faking it. Either the messages themselves should be idempotent, meaning they can be applied more than once without adverse effects, or we remove the need for idempotency through deduplication.
Honestly I don't get why this is "faking it" though. It seems like the author's definition of "exactly once" is so purist as to essentially be a strawman. This is "exactly once" in practice.
Like are there other people claiming that this purist version of exactly-once does exist?
In my experience, the purist version of "exactly-once" exists as a vague, wishy-washy mental model in the brains of developers who have never thought hard about this stuff[0]. Like, once you sketch out why idempotency is important and how to do it, folks seem to pick up on it pretty quickly, but not everyone has trained their intuition to where they automatically notice these sorts of failure modes.
[0] I don't mean this as a slight against those developers--the issues that arise from distributed systems are both myriad and subtle, and if you've spent your time learning how to make beautiful web pages or cool video games or efficient embedded systems, it seems reasonable to not know anything about the accursed problems of hypothetical Byzantine Generals. Or maybe you're fresh out of a bootcamp or an undergraduate program and haven't yet been trained to expect computers to always and constantly fail in every possible way.
It can get way harder when your initial design made incorrect assumptions about the delivery semantics you were using, so you didn't know you'd need it.
Edit for example:
Someone could have a low-latency problem that seems like it could be a fit for a streaming application. They could look at docs and see "ooh, with Flink I can do exactly-once writes to Kafka" in one place, and choose to use that. But if they don't dig deeply into what that means, they may miss the latency impacts of having to checkpoint every time to commit a set of writes to Kafka. And by the time they figure this out, managing both "low latency" and "exactly once" in the code they wrote might be a really hairy problem.
I've seen very few systems that have general idempotency baked in. Often it ends up being specific to the application. In some cases you can have simple solutions like upon crashing reload all of the state from an authoritative source. In some cases your messages result in simple idempotent operations such as "insert message with a unique ID" or "mark a message with a unique ID as read" but even then these are becoming quite related to business logic.
Basically idempotency is a powerful tool to create a solution but it is no silver bullet. That is why it is important to understand the underlying problem.
1. buy plane ticket 2. bring box to recipient 3. plug in Ethernet & send message
keep an eye out for our IPO
But that's not because you built a system that successfully delivers messages exactly once... you build a system that successfully processes messages exactly once, even if delivery occurs multiple times. The delivery still occurred multiple times. Even if your processing layer handled it, that may have other consequences worth understanding. Wrapping that up in a library may present a nice API for some programmer, but it doesn't solve the Byzantine General problem.
Whenever someone insists they can build Exactly Once with [mumble mumble mumble great tech here] I guarantee you there's a non-empty set of human readers coming away with the idea they can successfully create systems based on exactly-once delivery. After all, I built some code based on exactly-once delivery last night and it's working fine on my home ethernet even after I push billions of messages through it.
We're really better of pushing "There is no such thing as Exactly Once, and the way you deal with is [idempotence/id tracking/whatever]", not "Yes there is such a thing as Exactly Once delivery (see fine print about how I'm redefining this term)". The former produces more accurate models in human brains about what is going on and is more likely to be understood as a set of engineering tradeoffs. The latter seems to produce a lot of confusion and people not understanding that their "Exactly Once" solution isn't a magic total solution to the problem, but is in fact a particular point on the engineering tradeoff spectrum. In particular, the "exactly once" solutions can be the wrong choice for certain problems, like multiplayer game state updates, where it may be a lot more viable to think 1-or-0 and some timestamping and the ability to miss messages entirely and recover, rather than building an "exactly once" system.
I think the difference might be partly semantic. If processing at the messaging level is idempotent + at least once, then message delivery to the application level is exactly once. People mostly only care about the application level not the lower levels where they might just build on a library or system that handles that logic for them.
Alternatively we could come up with names for all the other combinations of delivery mechanism and handling mechanism, but since you can easily see we hit an NxM-type problem on that, this may well help elucidate why I think it's a bad idea to try to combine the two into one term. It visibly impairs people's ability to think about this topic clearly.
I'll grant that it matters if you're trying to debug some problem and trying to find at what layer it failed, but it's basically the same process you use to debug all of those other layers too, so I'm not sure why this layer deserves special consideration.
The problem with this is similar to the problems with two-phase commit in distributed databases: there are unavoidable failure cases. Most of the time it works just fine, but if you write your application to depend on this impossible feature, and it fails - which, given enough time, will certainly happen - then the cleaning up the mess can be much more effort (and have much wider business implications) than simply dealing with the undesirable behaviour of reality in the first place.
Or to put it another way: exactly once semantics can never be reliably extracted away from the application, so if you need it, it needs to be part of your application.
Hit an error, roll-back, side-affect can’t be rolled back. Retry - side-affect happens again.
Wouldn’t the general approach be to have unique message identifiers and queue side-affects? Maybe I’m missing lots of subtleties.
Idempotency (either via a token in the request, or another API to check the result of the previous request) is required to prevent duplicates. And this requires the third party service to support idempotency; there's nothing you can do on your side to enable it if their service doesn't support it.
If you guarantee "exactly once", you design your systems differently than "at least one with idempotence". A system designed for exactly once will be less complicated than a system designed for at least once + idempotence, which is why it is ideal but impossible.
Hey here's a solution to the halting problem – always assume yes, and then figure out the edge cases. How do you do that? Well that's on you, I did my job.
In a distributed system that needs exactly-once delivery, implementing perfect idempotence is equally impossible.
That's why if you search for exactly once delivery you'll see a bunch of products claiming to have it (e.g kafka).
idempotence + at least once
idempotence isn't necessarily commutative.
In active / active setups, there are other strategies such as partitioning and consensus.
There is a third option besides idempotency and eliminating side-effects: give each message a unique ID, use that to keep a record of which messages have been processed, and don't process the same message twice.
But I wouldn't get too worked up over it. The article is basically saying: You can't have exactly-once delivery (unless you take the necessary steps to ensure you have exactly once delivery).
Well, it's one way of implementing a kind of idempotency. But idempotency in general is more complicated than just deduping messages. For example, "Toggle the power state" is not idempotent because the result state depends on the initial state. You might think that "turn the power on" is idempotent, and by itself it is, but in conjunction with "turn the power off" it is not because the order in which they are processed matters. A truly idempotent message would be something like, "Insure that the number of on-off cycles at time T1 is N, at time T2 is N+1" etc.
Idempotency is in general more powerful and more complicated than deduping.
Here is a quote from the original article that supports my position:
"Therefore consumer applications will need to perform deduplication or handle incoming messages in an idempotent manner. ... The way we achieve exactly-once delivery in practice is by faking it. Either the messages themselves should be idempotent, meaning they can be applied more than once without adverse effects, or we remove the need for idempotency through deduplication."
Either you have a shared database and all consumers are local, in which case why are you passing messages at all, or you have a distributed system somehow. If you have a distributed system, you've got this problem, and going recursive probably won't help much.
Or I've missed something.
If it’s a notification displaying on a phone, the phone holds a DB that tracks what notifications have been shown.
You can’t reliably show a notification on exactly one of a user’s devices. But that’s no biggie; display it exactly once per device, remove it everywhere when acknowledged anywhere.
This is academic with enough abstraction, but when you're designing a system and implementing the work of processing a message it can be a pretty important distinction. Especially if you're past the point where the list of seen items becomes a chokepoint for synchronization between consumers.
> You can’t reliably show a notification on exactly one of a user’s devices. But that’s no biggie; display it exactly once per device, remove it everywhere when acknowledged anywhere.
Potentially messy. Now you have two distributed system messages - the initial notification and the ack.
Is it because the UUID is "only" random with some collision chance? (But then how come a sequential ID from your system wouldn't count?) Is it because my system needs to trust your system? (But why is trust a factor here?) Is it because my system needs a centralised database? (But our two systems are still distributed when you consider them together, right?) Is this a semantic argument over the meaning of "delivery", where I'm not allowed to impose requirements or check a database because the message has already been "delivered"? (But then why are we quibbling over semantics?) Is it because the message broker becomes stateful? (But why is that a constraint?)
I think from looking at the article that this is about delivery within a finite number of retries, but that seems like the kind of problem where in the real world we just ring each other when a message has been retried 50 times over the course of two days.
I think it's also worth noting that your suggested workaround has actually recreated the problem it's intending to solve. If you mark the UUID as "received" as soon as you get a message, how do you deal with duplicates if the processing fails? If you mark the UUID as "received" when you're done processing, how do you deal with the possibility that you'll receive the message multiple times? This can get hairy very quickly.
Every distributed systems engineer knows what "exactly-once delivery" is asking for, and in plain English it's valid to conflate the two, but for some reason the field has decided to treat the phrase as an annoying semantic pit trap for the unwary. Want to add an ID to your transactions to make them idempotent? Well, even though your transactions are now recorded exactly once, that wasn't technically delivery! Gotcha!
When you call something exactly-once, people who are perhaps not distributed systems engineers make the reasonable assumption that this means exactly what it says. They will engineer around this reasonable assumption based on a clear technical description and get something hilariously broken in non-obvious ways. This will have happened because jargon ("exactly-once delivery") has been confused for a technical description of a delivery system's properties.
Not everyone in this series of comments is a distributed systems engineer. Never mind everyone using a messaging system.
I’m arguing that exactly-once delivery is possible if the message receipt happens in a single place, but maybe others don’t see that as a distributed system at all.
Potentially messy. Now you have two distributed system messages - the initial notification and the ack.
Not a problem, both are idempotent! :)
I think that only holds if the sender is also in the same place. Otherwise there's the very real chance of a message getting lost, turning your system from exactly-once into at-most-once. At this point your sender, consumer, and messaging systems are one system, so it's probably reasonable to question if that's a distributed system.
No, they only need the ids of messages they have received.
After taking our customers through this same kind of apocalyptic rabbit hole conversation, they tend to agree with this architecture decision.
The cost of anticipating the .00001% that might never come is completely drowned out by the massive, daily 99%-certain headache that is managing a convoluted, multi-cloud cluster.
Many times the business owners will get the message and finally reveal that they have always had access to a completely ridiculous workaround involving literal paper & pen that is just as feasible in 2023 as it was in the 18th century.
It’s the resume-driven mid dev in the next office you’ve got to watch out for.
You seem to be confusing a system that produces bad results 1% of the time with a system that's down 1% of the time. If you can only write the first kind of non-distributed system, you're in for a bad trip if you try to write a distributed equivalent.
This led to an existential crisis because given the number of ports we open and the number of machines we run and the number of processes per machine, there must be over a 0.1% chance of any deployment blowing up this way. We do hundreds a year in prod and probably hundreds a month in preprod. We've been winning the lottery this whole time.
Throw enough events around and a one in a million corner case will happen every week, every day, twice a day, three times in a row. That gets old really really quickly.
I had to look this one up for a refresher, but 100% violently agree - Such behavior certainly warrants a bug submission.
Isn't this basically what every "reliable" transport (TCP, HTTP3, message queues...) does?
What determines which failure mode you get is whether the machine will failover to a machine that retries uncertain messages (giving you "at least once"), or it doesn't (giving you "at most once").
But, you say, why can't we have it failover to a machine that asks recipients what they have got and goes from there? Well we can, but the recipients don't know what messages are in the network still on their way to them.
But, you say, why not have the recipients disregard those inbound messages once they know about the replacement machine? Well you can do that, but now the *recipients* become machines whose job is to ensure the deduplication. And now *they* become the machine with a bad failure mode.
But, you say, does this not reduce the odds of failures? Why yes, it does. Which is why people do things like this. And there has to come a point where we accept SOME failure rate.
The alternative, well, read The Saddest Moment at https://scholar.harvard.edu/files/mickens/files/thesaddestmo... to see where madness leads.
What is the failure mode the recipients have here?
Frequently developers who implement these manage to do two bad things. First they introduced a lot of complex code which rarely triggers except in disaster, and so whose bugs tend to survive. And second, they manage to convince themselves that they have accomplished the impossible, and make reliability promises that other developers unwisely believe. Exactly how unreliable most systems were and how much the documentation couldn't be trusted was underappreciated until https://jepsen.io/ came along and started proving how bad most distributed software was.
Now it may seem bizarrely unlikely that you'll ever see this kind of situation. But failures often start from network congestion due to a packet storm. And the failures lead to chatty Byzantine fault tolerance protocols adding significant traffic. This causes cascading failures. And so a small, simple outage can escalate into a series of outages as ever more confused servers continue overwhelming the network with their futile attempts to discover what is supposed to be true. So complex combinations of failures occur together more often than most of us would expect.
Now you need a database. Do you also need exactly-once delivery to the database? Now the service is no longer stateless too, which means scalability is a problem. Maybe you decide to make it just an in-process cache for de-duping, but that needs expiring and now the semantics are exactly-once within a given time period, and not across service restarts.
We can definitely solve this with higher level constructs, but they're not free, and they can introduce the same issues themselves.
> Isn't this basically what every "reliable" transport (TCP, HTTP3, message queues...) does?
TCP does this, to solve retries at the TCP layer. HTTP3 does this to solve issues at the HTTP3 layer. Message queues might solve this for the message queue, depends. But none of these solve the product level, or the user experience level, or other higher levels where these issues still crop up. They're issues you have to solve at every layer in some way.
But yes, even if I needed some stuff to archieve it, that doesn't make it impossible as the OP claimed.
> But none of these solve the product level, or the user experience level, or other higher levels where these issues still crop up.
I don't understand this point. Do you have some examples?
How do you know what the last message you received is? You can crash in the middle of receiving a message. You can crash after you've written the ID to disk but before you've processed it, or you can crash after you've processed it but before you've written the ID to storage.
If you're a distribution box (which is quite, quite common in these message queue systems), you can get the message, and send it to a box that just powered off. You saw it, you recorded it, and now you're not sure if you forwarded it successfully or not.
Fun thing about power outages, not every box turns off at exactly the same nanosecond. PSUs are full of capacitors and inductors. Sometimes that's just enough to float through a brownout, too (and a bunch of machines booting can also cause a brownout)
This assumes you can generate monotonically increasing numbers. If you have many clients, now they all need to share a data source and may be performance bound by generating those numbers.
> Actually I don't even have to. It's probably a good thing to do for efficiency, but in principle I can just drop out-of-order messages and wait until they are redelivered, hopefully in the correct order
True (modulo first problem), but efficiency may be necessary here. With many clients, you may end up in a state where only a small fraction of messages get through successfully, and most traffic is unsuccessful, which is bad. This also makes performance commitments hard to maintain as it's perhaps just luck when a client manages to get a message through. Clients also now need more buffering, more state, etc.
>> But none of these solve the product level, or the user experience level, or other higher levels where these issues still crop up. >I don't understand this point. Do you have some examples?
Let's assume a simple client->server instant messaging app. As a user, if I send a message, I expect that to arrive exactly once. It's going over TCP which is "reliable", but it doesn't stop the HTTP request from failing and needing to be retried. It's using HTTP3, but that doesn't stop the server generating a 503 and needing to retry the POST request (or whatever). Maybe the server puts the message in a message queue, but that connection fails after sending a transaction commit, did it get committed?
Idempotency tokens or an equivalent mechanism do solve this, but there isn't one magic trick to solving it in some base layer technology like TCP, this needs to be solved again and again whenever you have distributed systems.
Also, this isn't just networking. Two processes on a single machine communicating via IPC may be effectively a distributed system! I've got some experience doing this on Android, and it's still hard.
You can assign an ID to your nodes and let them generate increasing numbers on their own. Node ID decides on a tie, and if one node sees a larger counter value appear, it adjusts its own counter so that it doesn't stay behind:
With this approach you'd still need to communicate the current clock number back to clients as otherwise one will get ahead and have all its traffic accepted, and others will fall behind and be unable to get traffic accepted. Even if an error causes a client to bump forwards to retry, by the time it has done that the number it is about to retry with may have been used.
Additionally, the aim is still to get exactly-once delivery, so clients need to be able to differentiate between an error caused by them reusing an ID that was rejected to enforce exactly-once delivery, and an error caused by another client getting that ID.
Basically, this issue is easy to solve with low traffic and reliable persistent storage everywhere, but hard to solve with high traffic, or the acceptance that all persistent storage brings additional reliability challenges.
Yes, with uniqueness constraint.
the service is no longer stateless too, which means scalability is a problem
Do you have a specific problem in mind?
Stateful services are far harder to scale than stateless ones. Typically a stateless service can be scaled out horizontally with relative ease, but when there's state storage involved this becomes harder. You can scale vertically, but only so far. You can scale horizontally, but then typically need to introduce some sort of sharding/consistent hashing/etc to ensure that the service has a consistent view of the world without needing to connect to every database instance.
The first time I encountered RabbitMQ it could only handle 60 messages per second with delivery guarantees. We already had our own bespoke system that used a database to handle a couple multiples of that. So we ended up limping along with what we had.
Also, isn't the assumption here that you will have a reliable connection to a shared DB?
You can have engineered solutions that is pragmatically close to deliver exactly once but it's not "pure" -- there are still scenarios, however unlikely, that it will fail.
The OP defines "distributed systems" like this:
> Web browser and server? Distributed. Server and database? Distributed. Server and message queue? Distributed.
By that definition, as soon as I have a server, a client and an unreliable connection between them, I have a distributed system. In that context, nothing stops me from counting IDs.
1. Some action with a side-effect (ex: update an entry on disk, send a message out to some third party, etc.). This might be a bank transaction, or a note saying "you gotta ship package X to person Y".
2. Some action to note that you've received ID X (ex: write to disk, send a message out to your DB, etc.)
How do you set up your server to deterministically do neither in event of a crash, and both in event that your code turns to completion?
So the action is some local action. If it’s purely digital, it’s fiddly but surely not impossible to ensure that the action and the record of the action either both take place or both don’t. It’s a database with a transaction log.
If the action is some irrevocable physical thing - remotely controlling a printer, say - you need to make a best effort to handle errors gracefully, sure.
I’d concede that it’s impossible to ensure that a document is printed out exactly once, say - maybe there’s a paper jam, and it’s debatable whether the jammed paper counts as a valid printout. But that’s not very surprising and I don’t think it tells you much. It’s mostly a problem for printer manufacturers rather than distributed system architects.
Edit: also, i should be fair and acknowledge that you're effectively describing idempotency (i'm guessing you already knew that ;P ), which the article's author eventually points out is a way to recover "exactly-once" semantics. The point, maybe, is that someone needs to explicitly do this somewhere; you can't really rely on your protocol to do it for you.
If your processing doesn’t have any external side effects (make external API calls, send emails, charge credit cards, etc) then one option is to put your message queue in a relational DB alongside all your other data. Then you can pull a message, process it, write the results to the DB, and mark the message as processed all inside one big transaction. But not many use-cases can fit these constraints and it also has low throughout.
These are impossibility results given various assumptions and requirements that may not hold in practice, or may be too restrictive. For one thing, i suppose we're pretty happy with a probabilistic solution as long as we can get the probability below an acceptable threshold.
If you get a letter with the same number you already read you don't even open it.
In Kafka this is also handled this way, events are numbered, and you request "latest" from the last one you processed.
In our eventstreaming it's also done like this, it may surprise you that Kafka is just an implementation of eventstreaming and not the same as.
With Kafka the offset of the consumers, or until which number it had already processed, used to be handled by the ZooKeeper but is migrated to the consumers.
There is on "exactly one consumer gets the message" done by the ZooKeeper, all consumers get all the messages from the topics they subscribed to. If you want exactly on exactly one consumer you should create different topics.
So not true.
So, in essence, you can never have "I will get this exactly once", and at best you can only ever have "I will have a plan for what to do if I get this more than once".
The question here would be who "you" are. Are "you" the low-level system processing raw messages, or are "you" a system on top of that?
The high-level system can rightfully claim to that it receives messages "exactly once"—from the low-level system.
I dont understand this at all. TCP has exactly once message delivery that the application layer is completely unaware of...
I guess hardly anyone does this though.
You need some processing for anything on the network, either you accept it or not. Yet, network protocols are described by the behavior they export to their consumers, not by their internals. Well, with the single exception of exactly-once delivery.
It absolutely can. That's the entire point of TCP.
If the network loses the last packet in a TCP stream, then goes down for an indefinite period, the two ends of the connection have irreconcilable differences of understanding about whether the entire transmission was received. As far as the recipient is concerned, they have an entire message that's fully ack'd and so they should process it. As far as the sender is concerned, they have no ack for the last chunk of data and must redeliver it. This is the crux of the Byzantine Generals problem.
TCP can solve problems that happen in its own domain, and give reliable in-order delivery once over an unreliable network up to a point. It's not able to provide exactly once semantics in all scenarios though. Because that's logically impossible.
Isn't that just shifting the goal a bit? Now the trick is "only process each number once" which seems to have its own transactional issues if you can crash between "taking action based on the message" and "recording that this number has been seen"? If you need non-idempotent actions wouldn't this still be a potential issue?
You Cannot Have Exactly-Once Delivery - https://news.ycombinator.com/item?id=9266725 - March 2015 (55 comments)
I suspect that quote will meet Resistance.
It was a system built by people that also didn't have distributed system experiences. It was not enjoyable at all, and at least once delivery was a consistent headache that required infrequent but time consuming remediation.
It's the case here, but isn't automatic.
Using a transaction to retrieve an item from the queue, and locking the row using "SELECT FOR UPDATE" and "SKIP LOCKED". Such that the row gets locked on read, and several workers can read from the table at the same time. Within the same transaction, other work is done, and everything gets committed to the database as a single atomic operation.
CockroachDB (a consensus/raft distributed database) recently added supported for SKIP LOCKED, but I still have yet to work on this idea.
Of course there are other concerns with using a database as a queue (mostly at high throughput) but for most cases it will work well.
I guess everyone has to make this mistake once in their career.
Funny enough, when I searched for “database as a queue”, my own comment from four years ago came up as the fourth result.
> The way we achieve exactly-once delivery in practice is by faking it. Either the messages themselves should be idempotent, meaning they can be applied more than once without adverse effects, or we remove the need for idempotency through deduplication.
Doesn't have to be a "true" exactly-once. Just practically so. It's like saying "humans can't actually fly... they are faking it by using machines".
--
[1] https://nats.io/
[2] https://docs.nats.io/using-nats/developer/develop_jetstream/...
It is a bit like saying that we can't have straight lines. Of course, if you zoom in far enough to see individual atoms, every physical surface will look jagged. But in practical terms we can have straight lines and surfaces to a good enough approximation. It means specifying what "straight" means and figuring out how to measure it and how to produce "straight" according to specification and measurements.
Engineering is about knowing and making tradeoffs. Every device we have ever created has to contend with limitations of physical existence. Engineering is about accomplishing goals in presence of those limitations.
A person who says "you can't deliver a message exactly once" clearly lives in an idealised, theoretical world. I would urge to leave your ivory tower for a second and see how engineers in real world accomplish what you say is not possible.
I get that this knowledge is useful -- but don't publish it as gospel. "You cannot have exactly-once delivery" is true, but not the same kind of truth as "you can't travel faster than light". No engineering can get you to travel faster than light. But engineering can get you as close to exactly-once delivery as you want to the point where the original statement stops being meaningful for real life problems.
Like someone else said, you can use at least once delivery and handle duplicate messages, but that's not quite the same as a distributed system guaranteeing that a message will be delivered exactly once.
Assuming you are talking about transferring between institutions, there is actually no single piece of software with this responsibility. The business processes are effectively what provide these guarantees (typically by way of another 3rd party).
In order to accomplish this, added latency (settlement time) is necessarily introduced into the process.
This isn't an opinion. This is a fact of distributed systems. An axiom, if you will.
so while this is a hugely important result, it doesn't stop us from building useful systems.
It’s just an interesting piece of theorising.
I think your comment (particularly your 4th paragraph about ivory towers etc) comes across as overly harsh and a little aggressive.
Assuming the network will recover in time may be a reasonable assumption sometimes, though.
This has nothing to do with CAP.
An axiom, by the very definition, cannot be proven.
Since the mathematical impossibility of exactly once delivery can be proven (there exists a proof), it means it is not an axiom.
Now, a mathematical proof is different from physical reality. In mathematics 10^(10^(10^10)) is different from infinity while in reality it is not. In mathematics you can halve distance between you and a point every second and you will never reach it while in physical world you will reach it in less than a minute. In mathematics you can win any lottery however low your win chance is by just trying it an infinite number of times. In real world you can't.
An anecdote (funny regardless of whether true) says that CIA cryptographers created once an unbreakable encryption scheme using one time pad. One time pad is mathematically proven to be impossible to break.
USSR cryptographers promptly learned to decipher all cryptograms.
Apparently, when combining plaintext and one time pad the machine would generate different voltages depending on the bit used in the pad, so rather than output 0V for false and 1V for true, the output became 0.9V or 1.1V for true and 0V or 0.1V for false depending on what bit was used in the pad.
This shows how idealised world is different from engineering reality and forgetting about it can lead to large errors in judgment.
The same kind of errors in judgment as people saying "you cannot have exactly-once delivery".
Programs work in real world and not in idealised mathematical space. In real world, exactly once delivery is a solved problem for any practical purpose. Which is evidenced by all those systems that actually do, in fact, provide exactly once delivery.
If what he really means is "guaranteed to be delivered exactly once", then, yeah, no, of course not, because if your network goes down, you can't send anything.
Over a network.
You can, of course, just have one big fat machine and have two nodes in the machine and send the packet from the one node to another in the one big giant machine. There's no "network" to go down, unless somebody slices into the machine with a red-hot katana.
You can also have a fully-connected topology where every machine is hard-wired into every other machine, and every machine will route for every other machine. There's no real "network" to go down; as long as there are 2 hosts that aren't dead, a message from one will get to another.
What I find amusing is that well before you get into truly hard distributed system problems, for a sufficiently complex implementation/application, your apps will be riddled with so many bugs and so many operational problems that the thing is going to fall over way before a well-built network goes down.
Can’t you de-dupe by attaching a UUID to messages?
> The way we achieve exactly-once delivery in practice is by faking it. Either the messages themselves should be idempotent, meaning they can be applied more than once without adverse effects, or we remove the need for idempotency through deduplication.
The author is wildly overestimating the complexity of the solution by ignoring the existence of a very simple but critical concept: UUIDs
If the client generates and assigns a UUID to each message it sends, the receiver can easily check if a specific message was already received before by comparing it against previously received UUIDs and can discard duplicates.
The title of the article is misleading because exactly-once delivery is in fact possible by deduplication. The deduplication can occur automatically (without having to know any context about the message beyond its UUID); it can occur BEFORE the delivery of the message to the consumer... So the title and entire premise of the article is incorrect.
What the author should really be saying is that more than one copy of the same message may occasionally pass through the network due to failures... But they are misleading in suggesting that it's not possible to deliver exactly one copy or that deduplication is a herculean task when in fact it's trivial.
It seems like you're conflating exactly-once delivery with exactly-once processing. Your proposed solution is only necessary _because_ we can't have exactly-once delivery as TFA states.
> If the client generates and assigns a UUID to each message it sends, the receiver can easily check if a specific message was already received before by comparing it against previously received UUIDs and can discard duplicates.
You quoted the part of the article that suggests exactly your solution: deduplication.
The author's definition of 'message delivery' is incorrect and so is the title which is clickbait.
I do not have such a use case in mind, but the description above does not sound terribly unrealistic to me.
So given those facts, your design of using UUIDs is absolutely not full proof if the database where you stored your UUID is not in a consistent state.
Here’s some more useful info for your perusal (https://ucare.cs.uchicago.edu/pdf/socc14-cbs.pdf)
“ Data consistency means that all nodes or replicas agree on the same value of a data (or eventually agree in the context of eventual consistency). In reality, there are several cases (5%) where data consistency is violated and users get stale data or the system’s behavior becomes erratic. The root causes mainly come from logic bugs in operational protocols (43%), data races (29%) and failure handling problems (10%).1 Be- low we expand these problems.”
If you assume that the sender cannot be trusted to not reuse a single UUID for multiple messages (e.g. to trick the receiver), you can still work around that limitation by computing and comparing message hashes (e.g. sha256) on the receiver side... Don't even need UUIDs. Every time you receive a message, you can hash it and check if you're received a message with this hash before (storing the hash of each message as it is received). You can use the hash in place of a UUID though in this case you probably need to add some index to each message to ensure that each hash is unique over time (since a message with the exact same payload will be counted as the same message, even if broadcast a long time apart).
A common Kafka approach is to partition by key, so that a given UUID will only be placed on one partition, and we're guaranteed that any further messages with that key will also be placed on that partition, so handled by the same consumer.
Then create a change-log topic that's co-partitioned by key with the input topic. And then funk around with partition assignment strategies so if consumer X is assigned partition 10 of the input topic, then it's also assigned partition 10 of the change-log topic.
Then add RocksDB as a fast state store on a persistent volume claim, as restoring state fresh from your change-log topic turns out to take about 7 minutes.
And then realise you've just reimplemented bits of Kafka Streams poorly.