Pathway Persistence Modes: RealtimeReplay, Batch, and SpeedrunReplay Use Cases

Pathway’s REALTIME_REPLAY respects original timestamps for event-time testing, BATCH loads snapshots instantly for analytics workloads, and SPEEDRUN_REPLAY processes data as fast as possible for CI validation—all configured via the PersistenceMode enum in pathway/internals/api.py.

Pathway (pathwaycom/pathway) provides a durable persistence layer that captures pipeline state and enables deterministic replay through the pw.PersistenceMode enumeration. These modes allow developers to reconstruct streaming behavior for debugging, batch analytics, or automated testing without maintaining live data sources.

RealtimeReplay: Event-Time Accurate Simulation

The pw.PersistenceMode.REALTIME_REPLAY mode replays persisted data as if it were arriving in real time, respecting the timestamps stored in the original snapshot.

Characteristics

  • Temporal fidelity: Downstream time-based operators such as windows, timers, and sessionizers see the same temporal ordering and pacing as during the original execution.
  • Event-time semantics: Uses the original event timestamps rather than processing time, ensuring deterministic behavior for event-time calculations.
  • Source simulation: Effectively emulates the original live data source without requiring active connections to external systems.

Use Cases

  • Testing event-time logic: Validate windowing aggregations or watermark behavior using historical data that exhibits the same timing characteristics as production.
  • Debugging temporal operators: Reproduce race conditions or timing-specific bugs that only appear with specific inter-arrival times.
  • Rebuilding pipelines: Reconstruct a live-like environment for downstream services when the original source is unavailable or expensive to maintain.

Batch: Instant Historical Analytics

The pw.PersistenceMode.BATCH mode loads the entire persisted snapshot in one go and processes it without any timing constraints, treating the data as a static dataset.

Characteristics

  • Immediate availability: All data is available instantly to downstream operators, eliminating latency between records.
  • Unbounded processing: No temporal barriers exist between early and late data, enabling global aggregations across the full snapshot.
  • Classic batch semantics: Transforms the streaming pipeline into a batch job while preserving the computation graph logic.

Use Cases

  • Data science workloads: Run heavy statistical aggregations, machine learning training, or feature engineering on historical streaming data.
  • Snapshot materialization: Generate comprehensive reports or materialized views from streaming pipelines for downstream batch consumers.
  • Backfilling analytics: Reprocess historical windows with updated business logic without simulating the original timing delays.

SpeedrunReplay: Maximum Velocity Validation

The pw.PersistenceMode.SPEEDRUN_REPLAY mode skips original timestamps and processes the snapshot as fast as possible, typically in a single thread, prioritizing throughput over timing accuracy.

Characteristics

  • Maximum throughput: Processes records at machine speed without artificial delays, completing replay orders of magnitude faster than real time.
  • Single-threaded execution: Typically operates in a simplified execution context that minimizes overhead for rapid validation.
  • Logical correctness focus: Verifies that operators compute correct results without testing timing-dependent behavior.

Use Cases

  • Automated testing: Execute unit tests and integration tests in CI/CD pipelines where latency is irrelevant but speed is critical.
  • Logic validation: Verify that UDFs, joins, and aggregations produce correct outputs without waiting for time-based triggers.
  • Rapid prototyping: Iterate quickly on pipeline logic using cached data without the overhead of real-time simulation.

Prerequisites: Storage Persistence Modes

Before utilizing any replay mode, you must first persist data using one of the storage modes defined in pathway/persistence/__init__.py. The pw.persistence.Config class controls what state is captured:

  • PERSISTING: Full persistence of tables, columns, and internal operator state. Defined as the default in pathway/persistence/__init__.py.
  • UDF_CACHING: Stores only the mapping from UDF input parameters to their results, avoiding expensive recomputation of functions like LLM calls or ML inference.
  • OPERATOR_PERSISTING: Stores only internal operator state (aggregations, joins, windows) without raw input data or UDF results, offering the most storage-efficient option for high-throughput pipelines.

Implementation Examples

The following examples demonstrate how to configure persistence backends and apply the three replay modes using pw.run().

RealtimeReplay for Event-Time Testing

import pathway as pw
from pathlib import Path

backend = pw.persistence.Backend.filesystem(Path("/tmp/pw_storage"))
persistence_cfg = pw.persistence.Config(
    backend=backend,
    persistence_mode=pw.PersistenceMode.PERSISTING,
)

# Initial run with full persistence

pw.run(
    my_graph_definition,
    persistence_config=persistence_cfg,
)

# Later: replay with original timing intact

pw.run(
    my_graph_definition,
    persistence_mode=pw.PersistenceMode.REALTIME_REPLAY,
)

Batch Mode for Analytics


# UDF caching example with batch replay

backend = pw.persistence.Backend.filesystem(Path("/tmp/pw_udf_cache"))
persistence_cfg = pw.persistence.Config(
    backend=backend,
    persistence_mode=pw.PersistenceMode.UDF_CACHING,
)

@pw.udf
def expensive_embed(text: str) -> list[float]:
    # Costly transformer call

    return model.encode(text)

table = pw.io.csv.read("input.csv")
embedded = table.select(vector=expensive_embed(table.text))

# Store expensive UDF results

pw.run(embedded, persistence_config=persistence_cfg)

# Reprocess everything instantly in batch mode

pw.run(embedded, persistence_mode=pw.PersistenceMode.BATCH)

SpeedrunReplay for CI Validation


# Example pattern from pathway/tests/test_io.py

backend = pw.persistence.Backend.filesystem(Path("/tmp/pw_ci"))
persistence_cfg = pw.persistence.Config(
    backend=backend,
    persistence_mode=pw.PersistenceMode.OPERATOR_PERSISTING,
)

# Production run with operator state saved

pw.run(production_graph, persistence_config=persistence_cfg)

# CI validation: verify logic instantly

pw.run(production_graph, persistence_mode=pw.PersistenceMode.SPEEDRUN_REPLAY)

Summary

  • RealtimeReplay processes persisted data with original timestamps intact, essential for testing event-time logic and windowing behavior in pathway/internals/api.py.
  • Batch loads entire snapshots immediately without timing constraints, enabling heavy analytics and data science workloads on historical stream data.
  • SpeedrunReplay maximizes throughput by ignoring timestamps, ideal for CI/CD pipelines and unit testing scenarios demonstrated in pathway/tests/test_io.py.
  • All replay modes require prior persistence using pw.persistence.Config with a backend (filesystem, S3, Azure) configured in pathway/persistence/__init__.py.

Frequently Asked Questions

What is the difference between storage modes and replay modes in Pathway?

Storage modes (PERSISTING, UDF_CACHING, OPERATOR_PERSISTING) determine what state gets written to the backend during the initial run, as configured via pw.persistence.Config. Replay modes (REALTIME_REPLAY, BATCH, SPEEDRUN_REPLAY) determine how that stored data flows through the pipeline during subsequent runs, specified via the persistence_mode parameter in pw.run().

Can I switch between Batch and RealtimeReplay on the same snapshot?

Yes. Once you have persisted data using any storage mode, you can replay that same snapshot using BATCH for analytics and later using REALTIME_REPLAY for timing-sensitive testing. The replay mode does not alter the stored snapshot; it only affects how the engine consumes the data.

Which persistence mode uses the least storage space?

OPERATOR_PERSISTING is the most storage-efficient mode, as it only persists the internal state of operators like aggregations and joins without storing raw input data or UDF results. According to the docstring in pathway/persistence/__init__.py, this mode performs persistence "only over the state of internal operators," minimizing disk or object storage usage.

How does SpeedrunReplay affect time-based operators?

SPEEDRUN_REPLAY processes records as fast as possible without respecting the original timestamps, meaning time-based operators such as tumbling windows or session windows may see all data as arriving simultaneously. This mode validates logical correctness but does not test temporal behavior or watermark handling.

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 →