Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

Streaming & Message Queues

Progress0/13
x

Batch versus Streaming

  • Batch versus streaming12m
  • Streaming architectures12m
  • Exactly-once semantics14m

Why Queues

  • Why a queue12m
  • Topics, partitions, keys14m

Offsets and Groups

  • Offsets and commits14m
  • Consumer groups14m

Delivery and Time

  • At-least-once14m
  • Late and out of order14m
  • Watermarks and Spark14m

Stream-batch

  • Stream-batch unification12m
  • Replay from offset12m

Capstone

  • Capstone: late messages16m
Back to track
  1. Learn
  2. Streaming & Message Queues
  3. Batch versus Streaming
  4. Exactly-once semantics

Lesson 3 of 13 · Theory first, then run it

Exactly-once semantics

streamingpythonbeginner14 min

Overview

At-most-once can drop events. At-least-once can duplicate them. Exactly-once is hard because the bookmark commit and the write are two steps. Idempotent consumers close the gap.

On this page7 sections›
  1. 1The idea
  2. 2Why this exists
  3. 3Picture this
  4. 4A small example
  5. 5Common beginner questions
  6. 6What comes next
  7. 7Practice

The idea

Delivery semantics describe what you can promise after a crash. At-most-once means an event might never be processed. At-least-once means it might be processed more than once. Exactly-once means it is processed one time, even if the consumer dies and restarts.

The consumer keeps a bookmark of how far it has read. In Kafka that bookmark is an offset. Commit the offset too early (before the warehouse write) and a crash loses the event: at-most-once. Commit after the write and a crash can replay the event: at-least-once. Exactly-once needs the write and the commit to succeed together, or an idempotent write that makes a replay harmless.

This lesson uses Python lists as the log and an integer as the bookmark. FakeKafkaTopic is a Python class of lists. No Kafka broker runs in this tab. Later lessons reuse the same three names with more machinery.

Why this exists

A payments consumer writes each paid order into a warehouse table, then marks the event read. If it writes ORD-1 and crashes before the bookmark moves, restart reads ORD-1 again and inserts a second row. Revenue doubles. That is at-least-once landing in a sink that is not idempotent.

Flip the order: mark ORD-1 read, then write. Crash before the write and ORD-1 never lands. Finance is missing a payment. Most data pipelines would rather risk a duplicate they can key-deduplicate than a silent gap. Exactly-once is the slogan people want. The implementation is the part that is hard.

Picture this

Mailing a letter and ticking it off a clipboard is two actions. If you tick first and the letter never goes, you lost it. If you mail first and forget to tick, you might mail it again. Exactly-once is a single receipt that covers both.

The crash window: write versus bookmark
Read eventWrite to warehouseCommit offset(bookmark)Crash hereduplicates or drops

Two steps that are not one transaction: warehouse write, then offset commit. A crash in between creates duplicates or gaps.

SemanticCommit bookmarkIf the process crashesTypical sink risk
At-most-onceBefore the writeEvent may never be writtenLost rows
At-least-onceAfter the writeEvent may be written twiceDuplicate rows
Exactly-onceWith the write, or write is idempotentRetry does not change the resultNeeds transactions or keys

A small example

Run the example below in this tab. Read the input, follow the code, then check the output matches what you expect.

Python
def at_most_once(log, offset, sink):
    """Mark read first. A crash before append loses the event."""
    event = log[offset]
    next_offset = offset + 1
    sink.append(event)  # if this never runs, the event is gone
    return next_offset

def at_least_once(log, offset, sink):
    """Write first. A crash before returning can deliver twice."""
    event = log[offset]
    sink.append(event)
    return offset + 1

log = [{"order_id": "ORD-1"}]
print("at-least-once next offset", at_least_once(log, 0, []))

An idempotent consumer keeps a set of keys it has already applied. A replay of ORD-1 hits the set and skips the second insert. You still run at-least-once delivery. The sink makes the second delivery a no-op, which is how many teams get once-in-the-warehouse results without warehouse transactions.

Python
def idempotent_apply(event, seen, sink):
    key = event["order_id"]
    if key in seen:
        return "skip"
    sink.append(event)
    seen.add(key)
    return "write"

seen = set()
sink = []
first = idempotent_apply({"order_id": "ORD-1"}, seen, sink)
replay = idempotent_apply({"order_id": "ORD-1"}, seen, sink)
print(first, replay, sink)

result = ["at-most-once", "at-least-once", "exactly-once"]
print(result)
  1. Decide whether a lost event or a duplicate event hurts more. Data pipelines usually pick duplicate-and-dedupe.
  2. Implement at-least-once: write, then commit the bookmark.
  3. Make the write idempotent on a business key (order_id), or use a transactional sink that commits offset and write together.
  4. Treat 'exactly-once' as an end-to-end claim. A consumer flag without an idempotent sink is still at-least-once.

Common beginner questions

Why is exactly-once hard if databases have transactions?

The log and the warehouse are often two systems. A transaction that covers both is special machinery (and slower). Across two commits, a crash can still split the work.

Is idempotent the same as exactly-once?

Not quite. Idempotent means doing the work twice leaves the same state. You may still process twice. The user-visible result matches once. That is the practical pattern for most warehouse loads.

Should I use at-most-once to avoid duplicates?

Only if losing events is acceptable (some metrics, some traces). Payments, orders, and balances are not that class of data.

Commit offset and write are two steps

If those two steps are not one transaction, you have a crash window. Pretending a config flag closed it is how duplicate revenue reaches a dashboard.

You will see this again

Later in this track, delivery lessons use the same three names with offsets on FakeKafkaTopic. The object is still a Python list. The broker is still not here.

What comes next

You have the vocabulary: batch versus stream, Lambda versus Kappa, and the three delivery promises. Next you will learn why a message queue exists at all: decouple, buffer, and replay, using FakeKafkaTopic as a list.

Practice

Run Sample to print the three semantics. Then complete the exercise: store ["at-most-once", "at-least-once", "exactly-once"] in result and print it.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Rate:
Was this useful?
Streaming architecturesWhy a queue