How Pathway Achieves Exactly-Once Processing Using Transaction-Based Semantics
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:+1for inserts and-1for deletes.
Together, these columns form a unique signature for every logical change set. In 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, 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 (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 (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:
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 (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
timeanddiffcolumns 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_onceandright_exactly_onceflags to guarantee deterministic results when inputs are transactionally unique. autocommit_duration_mscontrols 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.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →