No such thing as exactly-once delivery
blog.sequinstream.com
blog.sequinstream.com
This shows up explicitly in blockchain systems, where confidence improves as the number of confirmation cycles increases.
Processors, despite being a distributed system (well really every ASIC is), don't typically suffer from your classical distributed systems style issues because they have no failures, or rather they only have certain catastrophic failure modes that are low probability and well controlled for. Your typical distributed system to contrast has constant failures.
See [1]. This was a very real problem in early multiprocessor computer design.
Memory arbiters are where atomic operations on memory are actually implemented.
I don’t think there are any high bandwidth multi-drop busses anymore. Electrically multi-drop busses don’t make it easy to ensure the needed signal integrity for high speed busses.
Within a fully synchronous system, it's possible to avoid this, though. But most multi-processor systems are not synchronous.
Is this actually true? A modern AMD or Intel CPU doesn't have one quartz crystal per core, so therefore all per-core clocks have to be derived from one clock source, meaning the relationship between all the clock frequencies and phases is deterministic, even if not specced in the datasheet.
In practice the probability of that is just driven low enough that it's not an issue.
Processors suffer from the classical distributed system style issues just the same, the difference is that as you said they tend to be more reliable and the configuration doesn't change at runtime, but another key aspect is that a lot of latency tradeoffs aren't as lethal due to the short distances within a chip.
Accessing data that is kept in the cache of another processor only costs you 100ns at most and a lot of that overhead comes from time spent on the cache coherency protocol. Modern directory based cache coherency protocols are complicated beasts. Each processor core needs to be aware which processors have a copy of the cache line it is about to access, so that it can request exclusive write access to that cache line and mark every single one of those copies dirty upon modification.
The way the distributed system problems manifest themselves on the higher levels is that cache coherency by itself doesn't get you correct software automatically. You still have to use synchronization primitives or atomics on the software level, because your registers are essentially stale read replicas. Compare and swap on the processor level is no different than compare and swap on mongodb.
The users still need to interact with this wrapped system, and it's those interactions that can fail, thus preventing the overall exactly-once delivery.
To flip it around, if this wasn't an issue then what you're suggesting already exists in the form of ACID databases.
Which is the same problem, but at a different layer.
An analogy would be pizza delivery. You wouldn't call someone who delivered 100 identical pizzas to you even though you only ate one of them "exactly once delivery".
And, analogously, that does seem to be what Sequin's system provides.
If someone named Bob was standing outside your yard discarding the 99 pizzas, then Bob would not be able to guarantee that you actually get exactly one pizza. It's possible that the one pizza Bob tries to give you gets rained on before Bob can actually hand it to you.
Only you yourself can actually implement a means of discarding extraneous pizzas, any middlemen introduced into the chain just adds another layer of indirection that is susceptible to failure.
Of course the rain may flip bits in any system, but that would amount to destruction of you. You won't worry about pizza then, but your clone will spring up and get the pizza from Bob. And Bob himself is similarly protected from single-point-of-failure problems.
Of course the whole Bob subsystem may be destroyed, but that's already beyond the threshold of the problem we choose. Otherwise we'd have the external system to check if we'd actually got the pizza, and resend it again.
EDIT: that's probably the definition of "only you" - the Bob being that close. Agree then.
It's only in systems where failure is a possibility that exactly once delivery becomes impossible.
Yes, you can have exactly-once delivery - https://news.ycombinator.com/item?id=41598006 - Sept 2024 (133 comments)
Also. Others?
You cannot have exactly-once delivery (2015) - https://news.ycombinator.com/item?id=34986691 - March 2023 (217 comments)
Show HN: Exactly-once delivery counter using Watermill messaging library - https://news.ycombinator.com/item?id=26888303 - April 2021 (9 comments)
High-performance, exactly-once, failure-oblivious distributed programming (2018) - https://news.ycombinator.com/item?id=20562733 - July 2019 (26 comments)
Exactly-once Semantics: How Kafka Does it - https://news.ycombinator.com/item?id=14670801 - June 2017 (39 comments)
Delivering Billions of Messages Exactly Once - https://news.ycombinator.com/item?id=14664405 - June 2017 (133 comments)
You Cannot Have Exactly-Once Delivery - https://news.ycombinator.com/item?id=9266725 - March 2015 (55 comments)
Exactly-Once Messaging in Kafka - https://news.ycombinator.com/item?id=8880612 - Jan 2015 (8 comments)
Exactly-Once Delivery May Not Be What You Want - https://news.ycombinator.com/item?id=8614264 - Nov 2014 (10 comments)
The author should note that the homepage at https://sequinstream.com/ does in fact claim that Sequin supports "Exactly-once delivery" under the "Key features" heading.
> Exactly-once processing: Messages must be ack'd after they're delivered.
> For example, imagine when a worker processes a message, it performs
> a side effect like sending an email. It very well could receive the
> message, send the email, but then due to a network error fail to
> acknowledge the message.
Since we are now deeply in "explaining the joke"-territory: in this reply I am pretending to be a message queue implementation that never receives the ack from a receiver. Counting the vast number of upvotes on the comment, at least one other person was deligthed by this chain of replies with the subtext being apparent. I have also found your parent comment worthy of a sensible chuckle for the same reason.It was a good video until the RST.
Part of the reason I find distributed and concurrent computing so fascinating is that we lose so many assumptions about how computers are supposed to work; ordering guarantees, data actually arriving to destinations, fairness, etc.
While it is of course possible for data to get lost or mangled in a single-threaded, non-distributed application, it's rare enough to where you often don't need to plan much around it. Once you add concurrency, and especially if you make it distributed, you suddenly have to plan around the eventuality of the messages not arriving, or arriving in the wrong order, or arriving multiple times.
Personally, I've gotten into the habit of PlusCal and TLA+ for everything, and modeling in the assumption that stuff might never arrive or arrive multiple times, which can work as a very good sanity check (though it's of course not a silver bullet).
If a message can be guaranteed to be received one or more times but results in a single entry being inserted into a database (e.g. in an idempotent way using UUIDs to enforce uniqueness), then clearly, a consumer reading the records from that database would only receive the message exactly once. If you consider the intended recipient to be the process, not the host, then the definition of 'delivery' allows you to claim that exactly-once delivery is possible.
Of course this is a pragmatic, not mathematical argument. In theory you can't mathematically guarantee even at least once delivery because you can't guarantee that there's not going to be a Carrington event which will bring down the internet and drive humanity to extinction; thus preventing the consumer server from being rebooted.
So it's not just a pointless argument about the semantics of the term "delivery": the fact that no communication channel can have exactly once delivery means that these systems are much more difficult to implement. For example, you can't just chain together 2 "exactly once" delivery systems into some longer one with a stateless middle node: instead you need some concept of idempotency that spans both systems or the middle node itself needs to de-duplicate (statefully) the first link and forward messages on to the second link.
Same with two-phase commit, it can be implemented perfectly given certain reasonable constraints and assumptions. The consumer process could crash, then on restart, it would use timestamps to check which record it processed last and resume processing starting with the next unprocessed one. If the consumer uses input messages to produce output messages somewhere else, the consumer could also itself perform a two-phase commit on its own outputs to account for the possibility of itself crashing part-way through processing an input message. It can coordinate both inputs and outputs perfectly.
If we define 'processed' as fully committed, then messages ca be processed exactly once. If it can be processed once, it can be delivered exactly once if you accept the simplest definition of that word.
What is often suggested when people say "exactly once delivery is impossible" is that duplicates cannot be avoided and this is usually a cop out.
I don't see how. Let's take a simple example of a the receiving end of a messaging system, where the messaging system can do whatever it wants to enforce idempotency or whatever semantics it wants, and calls, in-process, some `handler()` method of the application.
The useful processing happens somewhere in `handler()`. No matter what point you identify as the point where processing happens inside handler() it won't have delivered/processed-once semantics: it will potentially be called multiple times. The fact that the messaging system will internally have a commit step after handler() completes which is idempotent is irrelevant: as a user of the system you care how many times handler() is called.
If we have a chain of senders and receivers and absolutely wanted to avoid double processing instead of also maintaining idempotence down the line as described above. It is also technically possible (but I guess not worth the complexity); we could make it so that the receiver would add a flag 'about_to_process' on each record just before it starts processing the record and then it would change it to 'committed' once it finishes. If the receiver crashes before completing processing a message, it would see the 'about_to_process' flag and in that case it could ask the next process/receiver in the sequence if they already received a packet with that UUID and only send the output for that message if they have not.
It's very ugly but it's possible.
That is, delivered exactly-once seems to be a relative statement where one party can state they did deliver exactly-once while another party can disagree, and they're both correct relative to their local frames of reference.
So, like how GR has no global conservation of energy, just local, could one say that distributed systems have no global conservation of messages, only local?
Or am I just late night rambling again...
There was some blog post last week that supposedly proved the impossibility of exactly once delivery wrong, but that was quite expectedly another attempt to rename the standard terms.
so what are we arguing about?
(Of course plenty of people on the Internet use different definitions to arrive at the opposite conclusion but they are all wrong and I am right.)
The similar question in TCP is what happens when the sender writes to their socket and then loses the connection. At this point the sender doesn't know whether the receiver has the data or not. Both are possible. For the sender to recover it needs to re-connect and re-send the data. Thus the data is potentially delivered more than once and the receiver can use various strategies to deal with that.
But sure, within the boundary of a single TCP connection data (by definition) is never delivered twice.
In a streaming app framework (like Flink or Beam, or Kafka Streams), I can get an experience of exactly once "processing" by reverting to checkpoints on failure and re-processing.
But what's the difference between doing this within the internal stores of a streaming app or coordinating a similar checkpointing/recovery/"exactly-once" mechanism with an external "delivery destination"?
It's important to note that by processing here we don't just mean any "computing" but the act of "committing the side effects of the computation"
Eg a streaming framework will let you tally an exactly once counter. You can flush the value of that tally to a data store. External observers will see an eventually consistent exactly-once-delivery result.
But a system is built by many components.
You can build a reliable system out of unreliable components.
You can build a transactional system out of non-transactional components.
However you pay a price for that, because you have to bear the consequences of this substrate in other layers, often in your business logic too.
What if you can't flush? You can't guarantee a flush: the sender or receiver can get powered off or otherwise get stuck for an indeterminate amount of time. That flush is the delivery. The data, is it stored in memory? Then it was lost when the process was killed.
A common streaming app will consume from a distributed log (Kafka, Kinesis) and only commit offsets if/when the data from the last offset to the next committed offset is fully processed in exactly-once semantics.
When delivering to the destination, the same story can hold. If the delivery destination returns a failure, then it may or may not have committed the write, but the streaming app will ensure either that the correct (exactly-once) value will be eventually flushed, or will retry computing the tally.
This will give you eventually consistent exactly-once semantics unless you allow the destination system (or any system) to remain in outage forever. But then of course if that's true, then you won't get any semantics, exactly-once, at-least-once, or otherwise.
Thus the very start of this thread: exactly-once-delivery and exactly-once-processing are different and it matters. If it didn't matter, there would be no offsets that kafka is tracking because everything would just work.
But if the processing contains for example sending an email then you need to apply said actions for all downstream actions, which is where exactly once as seen by the outside viewer gets hard.
No, as a vendor you're either deliberately lying, or completely incompentent. Let's insert the words they don't want to.
"We do exactly once good enough for your application!"
How could you possibly know that?
"We do exactly once good enough for all practical purposes!"
So are my purposes practical? Am I a true scotsman?
"We do exactly once good enough for applications we're good enough for!"
Finally, a truthful claim. I'm not holding my breath.
Edit: More directly address parent:
It's not possible to even guarentee at-least-once in a finite amount of time, let alone exactly-once. For example, a high-frequency trading system wants to deliver exactly one copy of that 10-million-share order. And do it in under a microsecond. If you've got a 1.5 microsecond round trip, it's not possible to even get confirmation an order was received in that time limit, much less attempt another delivery. Your options are a) send once and hope. b) send lots of copies, with some kind of unique identifier so the receiver can do process-at-most-once, and hope you don't wind up sending to two different receiving servers.
https://www.hydrogen18.com/blog/aws-s3-event-notifications-p...
From what I can tell my blog post eventually motivated them to change this behavior
This terminology confuses my face. Maybe we need a better name for "exactly-once processing"? When I say "processing" I'm thinking of the actions that my program is going to take on said message - the very processing that they're talking about can fail. Can we not just say 'the message is held until acknowledgement'?
Hence my complaint that making statements like "our system offers exactly once processing" can be completely correct - from the point of view of the messaging system. But it can be misleading to the casual reader in understanding how it fits into the broader solution, because the messaging system cannot guarantee that my application processes it exactly once. I'm saying the terminology doesn't feel right.
Exactly-once delivery means solving the Byzantine Generals problem, and since that's not possible, neither is exactly-once delivery.
Which isn't to say you can't get a good enough approximation, but you'll still only have an approximation.
I remember working on a system in AWS that hooked a SNS topic into a Data Firehose, which then delivered the data to S3. I was expecting rare duplicates and the system accounted for that.
After a week I audited the data and found that the duplicate rate varied between 5% and 8% per hour - much much higher than anticipated.
And this is purely within a managed integration between two managed components.
Optimistically assuming the message has not been delivered yet will sometimes reduce average latency (most databases), but for very expensive messages it is better to check first (rsync)
* X in transit to B
* A asks B: Have you received X? B says No.
* X gets to B.
* A sends X.
I think it's a useful excercise to attempt to disprove notions such as this. We shouldn't take all these "truths" for granted. But it does take some diving into it properly to appreciate the nuances and eventually get the underlying (as opposed to reasoning about it on a higher level, as we all are here).
- Within a transactional system you can have exactly once processing (as they call it here).
- Crossing the boundaries between two transactions requires idempotency IDs to deduplicate redeliveries and restore exactly-once inside the new transaction.
Email is an oft-cited example but is a transactional system that supports idempotency IDs (the Message-ID). Deliver the same email twice with the same message IDs and email systems are supposed to deduplicate. Gmail at least would do that, iirc.
What I've found is that usually what's happening is that the "debaters" assume different definitions of the thing they're debating about, which means they're really debating about two entirely different things, with the same name, and that's the real reason why they've come to different conclusions.
Therefore any adult arguing the "no free will" side is either
1. Disingenuous
2. Purposefully playing Candyland as an adult.
Either of which is so damaging to the Ethos that I lack any desire to continue the debate.
There's nothing disingenuous or moronic about believing that physical processes might be fundamentally deterministic.
My argument is that a game of Candyland is fundamentally deterministic in the same way.
And, unlike the game, it is impossible to hold the oracle's view (so even if freewill is illusory, it's an unbreakable illusion).
In a superdeterministic world, the presence of an illusion of free will is baked into everything humans do anyways. If the world is superdeterministic, but you can't actually predict what's about to happen with any additional certainty, then it doesn't really change anything anyway. So, there's no reason to change the way you argue based on the answer to the question in my opinion.
Of course, now I'm arguing about whether arguing about free will makes any sense, which is perhaps even more silly, but alas.
Same when it comes to arguments about whether or not our reality is a simulation, because, again, I'm not going to live any differently either way. My life is still important to me, even if it's just the result of an algorithm run on some computer built by an advanced civilization.
Ugh, this is just a matter of semantics. Boring. If the receiver only acts on messages once, and will always have the opportunity to act on every single message, then for all practical purposes, it's exactly-once delivery.
Is this insightful? Surprising? What's the point of the debate?
It’d be faster but we don’t care. It also participates in an XA transaction…
In 10 years of running it practically non stop up to volumes of 100kmsg/s with hundreds of network and disk issues, I’ve never ever seen even once a dupe message. Sure as the article is it could happen but in practice it doesn’t.
Try `kill -9`ing your consumers or unplugging a network cable and see what happens.
Also, AMQP is most certainly not exactly once by any measure at all. Do you have any consistency checks or anything to validate duplicates?
Either way, you’re getting duplicate deliveries if you’re using >1 machine. And if you’re not, what are you adding to the conversion?
Saying “Well I’ve driven 100k miles and never crashed” doesn’t add much.
Also, not sure how you’re consistency checking 100k messages per second over that many years.
# Receive new mail message, incoming id '1234'.
# No other writer knows about this message.
$ echo "hello world" > incoming/1.1234.txt
# If files stay in incoming/ too long, it means they probably failed writing,
# and are thus stale and can be removed.
# When the message is done being written, the writer moves it to a new state.
$ mv incoming/1.1234.txt new/1.1234.txt
# The writer can then report to the sender that the message was accepted.
# In mail, this is seen as optional, as the whole point is to deliver the message,
# not to make sure you inform the sender that you will deliver the message.
# Here's the exactly-once bit:
# If you need assurance of exactly-once delivery, split out the delivery portion
# from the confirmation portion.
#
# Have one actor talk to the sender, and another actor queue the message in another
# temporary queue, like 'accepted/'. Once the message is queued there, you tell the
# client it is accepted, wait for successful termination of the connection, and then
# move the message to the 'new/' queue.
#
# This way, if there is an error between accepting the message and confirming with
# the sender, you can detect that and leave the message in 'accepted/' where you can
# later decide what to do with it. (But it won't get picked up for processing)
# The message is now ready to be picked up for processing.
# To process the message, a processor picks up the new message and moves it to a new state.
#
# The processor adds an identifier '5678', so the processor can identify which file it's working on.
# It can have extra logic to re-attempt processing on this ID if it fails.
$ mv new/1.1234.txt process/1.1234.5678.txt
# If the processor dies or takes too long, you can identify that (stale file, processor tracking its work, etc)
# and this can be moved back to 'new/' for another attempt. Moving it back to 'new/' is still atomic so
# there is no danger of another processor continuing to work on it.
# Once the processor is done processing, it will move to the complete state.
$ mv process/1.1234.5678.txt done/1.txt
# If the processing file no longer exists, it means it was either already done,
# or previously died and was moved back to 'new/'.
This process all depends on a database [filesystem] using synchronous atomic operations [mv]. If your distributed system can't handle that, yeah, you're gonna have a hard time.* The human to whom the message is addressed reads it on their screen.
* The email is inserted into that person's inbox. It's possible something will happen to the message after than point, before the addressee actually reads it, but that isn't what you're talking about.
* The email has been accepted by a server operating as an endpoint for the addressee (their company's Exchange server), _and_ acknowledgement has been received by the sending server.
* The above, but no guarentees about the sending server getting the acknowledgement.
etc.
[Edit: also, what "guarentee" means to you. 100%, 99.99999% is good enough, and so on.]
It is of course "from the perspective of the host issuing the call", but that can be resolved by a network server by blocking operations on a given file when such an operation has begun. And of course the filesystem has to actually do the right thing. I'm no filesystem expert, but I would assume one way to do it is to write the block(s) that have the updated inode maps in one operation. A journal should help prevent corruption from making a mess of this. And of course disk / filesystem / OS tuning to further ensure data durability.
* Deliver-at-least-once: Ever had something hang for a very long time while trying to access a NFS mount that wasn't available? Say while booting. That's because NFS has a mode (hard mount) that tries really, really hard to guarentee delivery. In the bad old days you could literally wait forever. Nowadays things usually give up after a while, sacrificing delivery for some kind of functionality. You can never guarentee delivery within a finite time period. "But!" you say, "you can add as many redundant network cards, paths, and servers as you need to get as many 9's as you want in your delivery guarentee." That helps. But a) you still can't guarentee delivery (what if you've swapped network cards, and the new one isn't configuring properly?), and b) you now have the problem of
* Network partitioning. It is entirely possible, and happens in real life, that part of the network can't talk to the rest. So now you've got, say, 4 servers and 20 clients in one partition, and 3 servers and 100 clients in the other. What do you do? It's provably impossible [1] to guarentee all three features "consistincy" (no read gets an incorrect answer), "availability" (non-failing nodes continue to function), and "partition tolerance" (messages between nodes may be delayed for an arbitrary amount of time).
* Plus other stuff such as load balancing, consistent security contexts, backups, restores, adding and removing hardware on the fly, rolling updates, yada yada yada.
In many cases you can engineer a solution that's good enough. Not always; sometimes you just have to fork out the $$$ for monster monolithic servers. (Which, technically, are internally distributed redundant systems, but the SLA 9's can go way up if its all glued together in one box.)
Furthermore, there is no API to force 2 different directories to write to disk simultaneously in the same transaction. The best option is to fsync() both directories to disk. However, until both directories are confirmed to have been synced to disk, you run the risk of the file being in both directories during a well timed crash. Just to make this 100% clear: there are 0 guarantees that rename() has hit the disk once rename() has returned. rename() is just like every other filesystem operation that gets written back to disk at some later point in time after the syscall has returned (ignoring things like sync mount options).
FYI, I have worked on filesystems in Linux, and one of the applications I worked on had tests where we intentionally rebooted the system in the middle of a write heavy workload. Your view of the filesystem world is insufficiently nuanced.