At-Least-Once vs Exactly-Once Consistency Guarantees in Pathway: A Complete Guide

Pathway defaults to at-least-once delivery, which prevents data loss but may emit duplicate rows after recovery, while exactly-once semantics require either a graceful shutdown or the enterprise edition to ensure each input row produces exactly one output row even after crashes.

Pathway is a high-performance data processing framework that maintains fault tolerance through persistent state management and transactional batch processing. Understanding the distinction between at-least-once and exactly-once consistency guarantees is essential for building reliable streaming pipelines that handle failures without data corruption or duplication. This guide examines how Pathway implements both guarantees through its persistence subsystem and runtime configuration.

How Pathway Processes Data in Transactional Batches

Pathway’s execution engine processes incoming data in transactional batches—small groups of input rows that are processed atomically and persisted together. When a pipeline restarts after a crash or manual stop, the engine reads the persisted snapshot and the last committed offsets of input sources to resume processing.

The core persistence logic is controlled by the PersistenceMode enum defined in src/connectors/mod.rs (lines 35-44), which determines how snapshots and offsets are handled during replay and shutdown. This enum includes modes such as Batch, Persisting, and OperatorPersisting, each governing how the engine advances time and commits state during recovery.

At-Least-Once Guarantee: Default Behavior and Crash Recovery

At-least-once processing is the default guarantee provided by the free edition of Pathway and occurs whenever a program terminates abruptly (e.g., via kill -9 or power loss).

According to the persistence documentation in docs/2.developers/4.user-guide/60.deployment/55.persistence.md (line 91), the engine always re-processes the last uncommitted transactional batch after recovery. Because the snapshot records only the offsets that have been seen—not whether the batch was fully written to the downstream sink—this guarantees that no data is lost, but the batch may be emitted twice.

The README.md (line 99) explicitly states: "The free version of Pathway gives the 'at least once' consistency while the enterprise version provides the 'exactly once' consistency."

When running without specific persistence configuration, Pathway stores minimal metadata. If the process crashes while a batch is in-flight, that batch will be re-emitted upon restart, potentially creating duplicate rows in your output sink.

Exactly-Once Guarantee: Enterprise Features and Graceful Shutdowns

Exactly-once semantics ensure that each input row produces exactly one output row, even after pipeline restarts. This guarantee is available in the enterprise edition or achieved through graceful shutdowns in any edition.

During a graceful shutdown (calling pw.stop() after pw.run()), Pathway finishes the current transactional batch, writes the output atomically, and then advances the persisted offset. Because the batch is fully committed before the snapshot is taken, the next start sees the batch as already processed and does not re-emit it.

As documented in the persistence guide, "exactly-once semantic can be guaranteed" when the pipeline completes its current work before termination, preventing the duplication that occurs during crash recovery.

Configuring Persistence and Consistency in Code

Basic At-Least-Once Pipeline

The following example demonstrates the default at-least-once behavior without explicit persistence configuration:

import pathway as pw

class InputSchema(pw.Schema):
    word: str

words = pw.io.csv.read("inputs/", schema=InputSchema)
word_counts = words.groupby(words.word).reduce(pw.reducers.count())
pw.io.jsonlines.write(word_counts, "result.jsonlines")

pw.run()

Without persistence configuration, the pipeline stores minimal state. A crash during batch processing results in duplicate rows after restart.

Enabling Persistence (Still At-Least-Once)

To enable crash recovery while maintaining at-least-once guarantees:

persistence_backend = pw.persistence.Backend.filesystem("./state/")
persistence_cfg = pw.persistence.Config(persistence_backend)

pw.run(persistence_config=persistence_cfg)

The pipeline now recovers from crashes using the filesystem backend, but the last uncommitted batch may still be duplicated during recovery.

Achieving Exactly-Once via Graceful Shutdown

To prevent duplicates during planned restarts:

pw.run(persistence_config=persistence_cfg)

# ... processing occurs ...

pw.stop()  # Triggers graceful termination

During pw.stop(), the engine finalizes the current transactional batch, flushes the sink in src/python_api.rs, and advances persisted offsets. When restarted, the completed batch is not re-processed.

Advanced: PersistenceMode Configuration

For low-level control, the runtime uses the PersistenceMode enum from src/connectors/mod.rs:

from pathway import connectors

# Mode controlling batch replay behavior

mode = connectors.PersistenceMode.Batch

The Batch mode and related variants define how the engine advances logical time before replaying snapshots, which is fundamental to the exactly-once commit protocol implemented in src/persistence/config.rs.

Summary

  • At-least-once is the default guarantee in Pathway’s free edition, ensuring no data loss but allowing duplicate rows after crash recovery.
  • Exactly-once requires either the enterprise edition or a graceful shutdown via pw.stop() to atomically commit the final batch and prevent re-processing.
  • The persistence subsystem in src/connectors/mod.rs uses the PersistenceMode enum to control how transactional batches are snapshotted and replayed.
  • Configuring a persistence.Config enables state recovery, but only graceful termination or enterprise features upgrade the guarantee from at-least-once to exactly-once.

Frequently Asked Questions

Does Pathway's free edition support exactly-once processing?

No, the free edition of Pathway defaults to at-least-once consistency guarantees according to the README.md (line 99). While you can achieve exactly-once behavior in the free edition by using graceful shutdowns via pw.stop(), the enterprise edition provides additional mechanisms to enforce exactly-once semantics even during unexpected failures.

How do I ensure exactly-once semantics during a planned restart?

Call pw.stop() to trigger a graceful shutdown after pw.run(). This ensures the current transactional batch is fully committed to the downstream sink and the persisted offsets are advanced before termination. When the pipeline restarts, it recognizes the batch as completed and does not re-emit it, providing exactly-once guarantees for that batch.

What happens to my data if the Pathway process crashes unexpectedly?

If the process crashes while a batch is in-flight (e.g., from kill -9 or power loss), Pathway will re-process the last uncommitted transactional batch upon restart. As documented in 55.persistence.md, this ensures no data is lost, but may result in duplicate rows in your output sink since the engine cannot determine whether the previous attempt successfully wrote to the downstream system.

Where is the consistency guarantee logic implemented in the Pathway source code?

The guarantee logic spans multiple files: the distinction between free and enterprise features is documented in README.md (line 99); the persistence behavior is defined in docs/2.developers/4.user-guide/60.deployment/55.persistence.md (line 91); the runtime modes controlling batch handling are implemented in src/connectors/mod.rs (lines 35-44) via the PersistenceMode enum; and the configuration wiring is located in src/persistence/config.rs and src/python_api.rs.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →