4 ms·
In streaming apps, the data is not only stored in memory, and failures are inevitable. A common streaming app will consume from a distributed log (Kafka, Kines
by wellpast 2y ago
In streaming apps, the data is not only stored in memory, and failures are inevitable.
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.
- sethammons 2y agoand the only reason why that all works (kafka, kineis, etc) is _because_ they take care of things like "what if the power goes out before we can flush"? And the reason why they have to think of those things is because exactly-once-delivery doesn't exist. If it did, they wouldn't have to do anything. There would be no offsets needed. Offsets are needed because exactly-once-delivery doesn't exist and to achieve exactly-once-processing the extra computation is required. 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.