# How Pathway Achieves Exactly-Once Processing Using Transaction-Based Semantics

> Learn how Pathway implements exactly-once processing semantics with transaction-based commits and snapshot persistence for atomic data handling and reliable recovery. Discover the technical details.

- Repository: [Pathway/pathway](https://github.com/pathwaycom/pathway)
- Tags: internals
- Published: 2026-03-06

---

**Pathway guarantees exactly-once processing by treating every minibatch of data as a transaction marked with hidden `time` and `diff` columns, enabling atomic commits and deterministic recovery through snapshot-based persistence.**

Pathway is a Python data processing framework designed for streaming and batch analytics with strong consistency guarantees. According to the pathwaycom/pathway source code, the engine implements exactly-once processing semantics by embedding transactional metadata directly into every data minibatch. This design ensures that even after unexpected failures, downstream systems receive every record exactly once without duplicates or omissions.

## The Transactional Minibatch Model

Pathway treats every processing unit as a **transactional minibatch** that carries two hidden columns: `time` and `diff`. These columns are injected automatically by the runtime whenever a table is written to an external sink such as Postgres, MySQL, or MongoDB.

- **`time`** – Acts as a unique batch identifier that monotonically increases with each commit.
- **`diff`** – Indicates the operation type: `+1` for inserts and `-1` for deletes.

Together, these columns form a unique signature for every logical change set. In [`python/pathway/io/postgres/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/postgres/__init__.py) (lines 38-45), the write operation automatically appends these columns to the output schema, ensuring that every commit to the external system is tagged with its transactional context.

## Snapshot-Based Persistence for Exactly-Once Recovery

When a Pathway program runs with persistence enabled, the engine periodically **snapshots** the state of all input tables to durable storage. This snapshot includes the hidden `time` and `diff` columns, preserving the complete transactional context.

On restart after a crash, the engine reloads the last consistent snapshot and resumes processing from the next `time` value. Because every minibatch has a distinct `time` identifier, Pathway can detect whether a specific batch was already applied to the sink and skip it if necessary. This mechanism, implemented in [`python/pathway/internals/table.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/table.py), ensures **exactly-once delivery** during graceful shutdowns and failure recovery.

## Window Operators with Exactly-Once Behavior

For temporal operations, Pathway provides `pw.temporal.exactly_once_behavior()` to prevent duplicate window emissions. This function returns an `ExactlyOnceBehavior` object that instructs window operators to emit **only one output** when a window closes.

Internally, the operator waits until the minibatch `time` exceeds the window’s end timestamp (plus any configured shift) and then freezes the result. Any subsequently arriving late events are discarded rather than triggering additional emissions. The implementation in [`python/pathway/stdlib/temporal/temporal_behavior.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/stdlib/temporal/temporal_behavior.py) (lines 83-99) ensures that tumbling or sliding windows produce deterministic, singular results even under reprocessing scenarios.

## Join-Level Exactly-Once Guarantees

Pathway extends exactly-once semantics to join operations through the `left_exactly_once` and `right_exactly_once` flags available in many join operators. When both input sides are marked as exactly-once, the join engine assumes each input row appears in **only one** transaction.

This assumption allows the join to emit deterministic results without generating duplicate rows when the job restarts. The logic is implemented in [`python/pathway/internals/joins.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/joins.py) (lines 143-165), where the engine optimizes the join execution based on these transactional guarantees.

## Transaction Grouping and Commit Control

The runtime groups multiple logical operations into a single transaction before flushing to the sink using the `autocommit_duration_ms` parameter. This configuration balances throughput with durability by reducing the frequency of intermediate snapshots while preserving the atomicity of each minibatch.

By tuning `autocommit_duration_ms`, you control the trade-off between latency and exactly-once safety. Shorter intervals minimize reprocessing scope during failures, while longer intervals improve throughput by batching more operations into a single transactional commit.

## End-to-End Exactly-Once Processing Example

The following example demonstrates a complete pipeline that reads from Kafka, applies windowed aggregations with exactly-once semantics, and writes to Postgres with transactional guarantees:

```python
import pathway as pw

# 1️⃣ Define a schema (no need to declare time/diff – they are added automatically)

class InputSchema(pw.Schema):
    user_id: int
    event_ts: int
    value: float

# 2️⃣ Read a Kafka topic (each incoming record is placed into a transaction)

source = pw.io.kafka.simple_read(
    "kafka:9092",
    topics=["events"],
    format="json",
    schema=InputSchema,
    autocommit_duration_ms=50,      # group records into minibatches

)

# 3️⃣ Window with exactly‑once semantics

windowed = source.window_by(
    pw.this.event_ts,
    interval=pw.timedelta.minutes(5),
    behavior=pw.temporal.exactly_once_behavior(),   # emit only once per window

).reduce(
    total=pw.reducers.sum(pw.this.value)
)

# 4️⃣ Write to Postgres – time/diff columns are created automatically

pw.io.postgres.write(
    table=windowed,
    postgres_settings={
        "host": "localhost",
        "port": "5432",
        "dbname": "analytics",
        "user": "app",
        "password": "secret",
    },
    table_name="metrics_5m",
    output_table_type="snapshot",   # snapshot stores only the latest window result

)

# 5️⃣ Enable persistence so that a crash can be recovered exactly‑once

pw.run(persistence_dir="/tmp/pathway_state")

```

This pipeline uses [`python/pathway/tests/temporal/test_windows_stream.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/tests/temporal/test_windows_stream.py) (lines 380-392) to validate the exactly-once behavior in practice, ensuring that window aggregations and sink writes remain consistent across restarts.

## Summary

- **Transactional minibatches** use hidden `time` and `diff` columns to uniquely identify atomic change sets in every output commit.
- **Snapshot-based persistence** captures these transactional columns, enabling the engine to resume from the last consistent state and skip already-processed batches.
- The **`exactly_once_behavior()`** API ensures window operators emit singular results by freezing outputs once the minibatch time exceeds the window boundary.
- **Join operators** accept `left_exactly_once` and `right_exactly_once` flags to guarantee deterministic results when inputs are transactionally unique.
- **`autocommit_duration_ms`** controls the grouping of operations into atomic transactions, balancing exactly-once safety with processing throughput.

## Frequently Asked Questions

### How does Pathway prevent duplicate records when writing to external databases?

Pathway prevents duplicates by tagging every minibatch with unique `time` and `diff` columns that act as a transactional identifier. When writing to Postgres or other supported sinks, these columns allow the engine to detect whether a specific batch was already committed. If a job restarts after a failure, the snapshot recovery mechanism resumes processing from the next uncommitted `time` value, ensuring previously applied batches are skipped.

### What happens if a Pathway job crashes during window processing?

If a job crashes during window processing, the persistence layer restores the last saved snapshot, including the hidden transactional columns. The engine then reprocesses data starting from the next minibatch `time`. For windows configured with `exactly_once_behavior()`, the operator checks the restored `time` values to determine if the window already closed and emitted its result, preventing duplicate emissions during recovery.

### Do I need to manually define the time and diff columns in my schema?

No, you do not need to manually declare the `time` and `diff` columns. Pathway adds these automatically when tables are written to external sinks or when persistence is enabled. These hidden columns are managed entirely by the runtime infrastructure, allowing your application code to focus on business logic while the engine handles transactional metadata transparently.

### When should I use exactly_once_behavior() versus standard window behavior?

Use `exactly_once_behavior()` when you require that each window emits exactly one final result that never changes, even if late data arrives after a crash recovery. Standard window behavior may emit corrections or updates as new data arrives, which is suitable for scenarios requiring real-time responsiveness over strict exactly-once guarantees. Choose exactly-once behavior for downstream systems that cannot handle retractions or duplicate window results.