What Persistence Backends Does Pathway Support for State Recovery and Fault Tolerance?
Pathway supports four persistence backends for fault-tolerant state recovery: local Filesystem, Amazon S3 (and MinIO-compatible services), Azure Blob Storage, and an in-memory Mock backend for testing.
Pathway, the open-source stream processing framework maintained by pathwaycom/pathway, provides a configurable checkpointing layer that stores operator state, frontiers, and metadata in a pluggable key-value storage layer. This architecture enables exactly-once semantics and automatic recovery from failures across both single-node and distributed deployments.
The Four Persistence Backends
The persistence layer in src/persistence/backends/ implements four concrete storage backends, all adhering to the PersistenceBackend trait defined in src/persistence/backends/mod.rs.
Filesystem Backend
The Filesystem backend stores snapshots and metadata in a local directory, making it suitable for development, single-node runs, or environments with shared network drives. Internally, Pathway uses FilesystemKVStorage (defined in src/persistence/backends/file.rs) to map key-value operations to files on disk. This backend requires no external dependencies and provides the lowest latency for local recovery.
Amazon S3 and Compatible Storage
The S3 backend writes snapshots as objects to Amazon S3 or compatible services like MinIO. Implemented in src/persistence/backends/s3.rs via S3KVStorage, this backend guarantees durability across clusters and supports multi-node deployments. Pathway organizes snapshot data under a configurable prefix within the bucket, separating metadata from per-operator state files.
Azure Blob Storage
For Azure-native deployments, Pathway provides the Azure backend (AzureKVStorage in src/persistence/backends/azure.rs). This backend offers identical semantics to S3 but targets Azure Blob Storage accounts, integrating with Azure's identity and networking models for enterprise cloud environments.
Mock Backend for Testing
The Mock backend provides an in-memory stub (MockKVStorage in src/persistence/backends/mock.rs) used exclusively for unit tests or when users explicitly require a no-op persistence layer. This backend does not survive process restarts but eliminates I/O overhead during testing.
Technical Architecture: The PersistenceBackend Trait
All four backends implement the PersistenceBackend trait, which defines the core key-value operations required by the checkpointing system: list_keys, get_value, put_value, and remove_key.
The enum PersistentStorageConfig (defined in src/persistence/config.rs at lines 46-59) selects the concrete implementation at runtime. When a pipeline starts, Pathway creates a PersistenceManagerOuterConfig containing the chosen backend, snapshot interval, and access flags. Internally, PersistenceManagerConfig transforms this configuration into three distinct storage objects:
create_metadata_storage– A single KV store per worker for pipeline metadata.create_snapshot_writerandcreate_snapshot_readers– Per-operator KV stores for state and offset snapshots.create_cached_object_storage– An optional local cache for large objects when using remote backends.
Pathway automatically adjusts the snapshot merging interval based on whether the backend is local or remote. Remote backends trigger less frequent merging through adjusted_snapshot_merging_interval to minimize API call costs and latency.
Configuring Persistence in Python
Pathway exposes backend selection through pw.persistence.Backend methods in the Python API. The bridge function construct_persistent_storage_config in src/python_api.rs maps string identifiers ("fs", "s3", "azure", "mock") to the appropriate PersistentStorageConfig variants.
Filesystem Configuration
import pathway as pw
backend = pw.persistence.Backend.filesystem("./Pathway-Cache")
config = pw.PersistenceConfig(
snapshot_interval=60000, # 60 seconds
backend=backend,
snapshot_access=pw.persistence.SnapshotAccess.All,
persistence_mode=pw.persistence.PersistenceMode.OperatorPersisting,
)
S3 or MinIO Configuration
import pathway as pw
from pathway.io.s3 import AwsS3Settings
s3_settings = AwsS3Settings(
bucket_name="my-bucket",
region="eu-central-1",
access_key="****",
secret_access_key="****",
)
backend = pw.persistence.Backend.s3(s3_settings, "pipeline-state/")
config = pw.PersistenceConfig(
snapshot_interval=120000,
backend=backend,
snapshot_access=pw.persistence.SnapshotAccess.ReadWrite,
persistence_mode=pw.persistence.PersistenceMode.OperatorPersisting,
)
Configuring Persistence in Rust
When using Pathway's Rust API directly, instantiate PersistentStorageConfig variants and wrap them in PersistenceManagerOuterConfig:
use pathway_engine::persistence::config::{
PersistentStorageConfig, PersistenceManagerOuterConfig, SnapshotAccess, PersistenceMode
};
use std::path::PathBuf;
use std::time::Duration;
let fs_cfg = PersistentStorageConfig::Filesystem(PathBuf::from("/tmp/pw-state"));
let outer = PersistenceManagerOuterConfig::new(
Duration::from_secs(60), // snapshot interval
fs_cfg,
SnapshotAccess::ReadWrite,
PersistenceMode::OperatorPersisting,
true, // enable compression
false, // disable async snapshots
Duration::from_secs(30), // merging interval
);
Summary
- Pathway provides four persistence backends: Filesystem, Amazon S3 (MinIO-compatible), Azure Blob Storage, and Mock.
- All backends implement the
PersistenceBackendtrait with standard KV operations (list_keys,get_value,put_value,remove_key). - The
PersistentStorageConfigenum insrc/persistence/config.rscontrols runtime backend selection. - Python users configure backends via
pw.persistence.Backend.filesystem(),Backend.s3(), orBackend.azure(). - Remote backends automatically trigger less frequent snapshot merging to optimize API costs.
- Core implementations reside in
src/persistence/backends/file.rs,s3.rs,azure.rs, andmock.rs.
Frequently Asked Questions
How do I choose between filesystem and cloud storage for Pathway persistence?
Use the Filesystem backend for local development, single-node deployments, or when operating in air-gapped environments with shared network storage. Choose S3 or Azure for distributed clusters requiring durability across machine failures or when running in Kubernetes across multiple availability zones. Cloud backends ensure state survives node termination but introduce network latency during snapshot operations.
Can I use MinIO or other S3-compatible storage with Pathway?
Yes. The S3 backend uses standard S3 API calls implemented in src/persistence/backends/s3.rs, making it compatible with MinIO, Ceph, and other S3-compatible object stores. Configure the endpoint URL and credentials through AwsS3Settings when initializing pw.persistence.Backend.s3().
What is the Mock backend used for in Pathway?
The Mock backend provides an in-memory hash map implementation of PersistenceBackend for unit testing and CI pipelines. It eliminates disk I/O and cloud dependencies, allowing tests to verify checkpointing logic without external infrastructure. Never use Mock in production, as all state disappears when the process exits.
How does Pathway handle snapshot intervals with different persistence backends?
Pathway adjusts the snapshot merging behavior based on backend latency. For remote backends like S3 and Azure, the system increases the adjusted_snapshot_merging_interval to reduce the number of ListObjects and PutObject API calls, lowering costs and preventing throttling. Local filesystem backends use more aggressive merging intervals since disk I/O is cheaper and faster.
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 →