# How Pathway Manages Backpressure and Optimizes Memory Usage in Streaming Dataflows

> Pathway optimizes streaming dataflows by managing backpressure and memory usage with a configurable max_backlog_size, preventing exhaustion and bounding RAM to specified capacity.

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

---

**Pathway prevents memory exhaustion in streaming dataflows by enforcing a configurable `max_backlog_size` limit that pauses input connectors when the number of in-flight rows exceeds the threshold, directly bounding RAM usage to the specified backlog capacity.**

Pathway is an open-source streaming data processing framework that combines a Rust execution engine with Python APIs to handle high-throughput dataflows. Managing backpressure and memory optimization is critical for production streaming pipelines, especially when dealing with bursty input sources like Kafka or filesystem directories. According to the pathwaycom/pathway source code, the engine implements a precise backpressure mechanism through the `max_backlog_size` parameter, which controls exactly how many rows can reside in memory at any given moment.

## Architecture Overview: Rust Core with Connector Isolation

Pathway’s architecture separates the high-performance Rust execution core from data ingestion connectors that run in independent Python or Rust threads. This design allows the engine to maintain precise control over the **in-flight data**—rows present in internal mini-batches that have been read but not yet fully processed. By isolating connectors in separate threads, Pathway can pause ingestion without blocking the computational engine, creating a natural boundary for backpressure implementation.

## The Backpressure Mechanism

Pathway implements backpressure through a cooperative protocol between input connectors and the Rust execution engine. When memory pressure rises, the system throttles data ingestion rather than accumulating unbounded queues.

### The max_backlog_size Parameter

Every input connector in Pathway—including filesystem, Kafka, S3, and Python subjects—accepts an optional **`max_backlog_size`** argument in `DataSourceOptions`. This integer specifies the maximum number of entries (rows or messages) that may be queued for processing at any moment.

When the aggregate size of all mini-batches currently held by the engine reaches this limit, the connector’s read loop is automatically paused. Reading resumes only after at least one batch has completely exited the engine—meaning all its outputs have been produced and the rows are no longer needed for computation.

### Engine-Side Enforcement

While the engine processes mini-batches, it continuously tracks the **aggregate number of rows** in all active batches. If admitting a new batch would exceed the configured `max_backlog_size`, the Rust side signals the connector thread to stop reading. This prevents uncontrolled memory growth during input spikes and ensures that RAM usage remains proportional to the configured backlog capacity rather than the total input volume.

### Connector Implementation Details

The `max_backlog_size` value propagates through the connector hierarchy:

- **Definition**: Stored in `DataSourceOptions` containing the `max_backlog_size: int | None` field in [`python/pathway/internals/datasource.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/datasource.py) at lines 22-23.
- **Propagation**: When a connector is built, the option transfers to a `GenericDataSource` via `_create_python_datasource` in [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py) at lines 31-38.
- **Python Subjects**: For Python subjects, the buffer is a `queue.Queue`. If `max_backlog_size` is supplied, `_set_max_backlog_size` replaces the default unlimited queue with a bounded one (`Queue(max_backlog_size)`) at lines 310-315 in [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py).

If `max_backlog_size` is omitted, the limit defaults to `None`, allowing the connector to read without restriction. This is convenient for low-volume workloads but risks unbounded memory usage during large bursts.

## Memory Optimization Strategies

Pathway’s memory optimization is not merely about limiting queues; it is integrated with the batch processing lifecycle to minimize resident data.

### Bounded Mini-Batches

Data is grouped into time-based mini-batches controlled by `autocommit_duration_ms`. Only whole batches are retained while they are needed for downstream operators. By combining `autocommit_duration_ms` with `max_backlog_size`, developers create a two-dimensional bound: limits on both the time window and the row count of in-flight data.

### Automatic Batch Eviction

Once every output of a batch has been emitted, the engine drops the batch’s rows, immediately freeing space for new data. This automatic eviction works in concert with the backpressure system: when the backlog limit would be exceeded, the pause stops additional rows from entering the system, preventing the queue from growing beyond the user-specified bound while the engine clears existing batches.

Together, these mechanisms guarantee that a streaming pipeline cannot consume more RAM than the developer explicitly allows through the `max_backlog_size` parameter.

## Practical Implementation Examples

The following examples demonstrate how to configure backpressure limits in real-world scenarios.

### Filesystem Source with Backpressure

Read a directory of JSON files while limiting memory usage to 500 rows:

```python
import pathway as pw

table = pw.io.fs.read(
    "data/incoming/",
    format="json",
    mode="streaming",
    max_backlog_size=500,
    autocommit_duration_ms=20,
)

out = table.select(lambda r: {"id": r.id, "value": r.value})
pw.io.csv.write(out, "data/out/processed.csv")
pw.run()

```

### Custom Python Subject with Bounded Queue

Create a custom connector that throttles when the backlog reaches 2,000 entries:

```python
import pathway as pw

class MySubject(pw.io.python.ConnectorSubject):
    def run(self):
        for i in range(10_000):
            self.next(id=i, payload=str(i))

subject = MySubject()
table = pw.io.python.read(
    subject,
    schema=pw.Schema,
    max_backlog_size=2_000,
    autocommit_duration_ms=10,
)

count = pw.run(pw.debug.compute_and_print(table.count()))

```

### Kafka Source with Memory Limits

Consume from Kafka while preventing unbounded growth during traffic spikes:

```python
import pathway as pw

kafka_cfg = {
    "bootstrap.servers": "localhost:9092",
    "group.id": "demo",
    "auto.offset.reset": "earliest",
}

table = pw.io.kafka.read(
    topics=["events"],
    config=kafka_cfg,
    format="json",
    max_backlog_size=10_000,
    autocommit_duration_ms=5,
)

filtered = table.where(lambda r: r.type == "click")
agg = filtered.groupby(lambda r: r.user_id).agg(count=pw.count())
pw.io.jsonlines.write(agg, "out/agg.jsonl")
pw.run()

```

## Key Source Locations

The backpressure implementation spans the connector interface and engine coordination. These files contain the definitive implementation details:

- **[`python/pathway/internals/datasource.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/datasource.py)** (lines 22-23): Defines `DataSourceOptions` containing the `max_backlog_size: int | None` field.

- **[`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py)** (lines 31-38): Implements `_create_python_datasource` which propagates `max_backlog_size` to `GenericDataSource`.

- **[`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py)** (lines 310-315): Contains `_set_max_backlog_size` which instantiates a bounded `Queue(max_backlog_size)` for Python subjects.

- **[`docs/2.developers/4.user-guide/80.advanced/50.how_pathway_connectors_work.md`](https://github.com/pathwaycom/pathway/blob/main/docs/2.developers/4.user-guide/80.advanced/50.how_pathway_connectors_work.md)** (lines 86-98): Documents the conceptual pause/resume logic and backpressure control semantics.

- **[`python/pathway/tests/test_io.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/tests/test_io.py)** (lines 5183-5250): Houses `test_backpressure_management_*` scenarios that verify atomicity, list-of-objects support, and rewind behavior.

## Summary

Pathway prevents memory exhaustion in streaming pipelines through a cooperative backpressure mechanism that bounds in-flight data:

- **Configurable Limits**: The `max_backlog_size` parameter sets a hard cap on queued rows across all input connectors, including filesystem, Kafka, S3, and Python subjects.
- **Engine Coordination**: The Rust core tracks aggregate row counts and signals connector threads to pause reading when limits are reached, resuming only after batches complete processing.
- **Memory Safety**: By combining `max_backlog_size` with `autocommit_duration_ms`, developers create a two-dimensional bound on both row count and time windows, ensuring RAM usage remains proportional to configured capacity rather than input volume.
- **Implementation**: The mechanism is implemented in `DataSourceOptions`, enforced through bounded `queue.Queue` instances for Python connectors, and coordinated by the Rust engine as documented in the advanced connector guide.

## Frequently Asked Questions

### How does max_backlog_size prevent out-of-memory errors?

The `max_backlog_size` parameter places a hard limit on the number of rows that can be held in memory simultaneously across all active mini-batches. When the engine detects that accepting a new batch would exceed this limit, it signals the connector to pause reading. This throttling ensures that memory consumption cannot grow beyond the configured bound, even during traffic spikes or when downstream operators are slow.

### What happens if I don't specify max_backlog_size?

If omitted, `max_backlog_size` defaults to `None`, which configures the connector to use an unbounded queue. While this eliminates throttling overhead for low-volume workloads, it removes the safety guarantee against memory exhaustion. During input bursts, the engine will accumulate all incoming rows in mini-batches until system memory is depleted, potentially causing the process to crash or be killed by the OS.

### Can I use backpressure with custom Python connectors?

Yes. Custom connectors built using `pw.io.python.ConnectorSubject` fully support backpressure through the `max_backlog_size` parameter in `pw.io.python.read()`. When specified, Pathway invokes `_set_max_backlog_size` to replace the default unlimited `queue.Queue` with a bounded version (`Queue(max_backlog_size)`). This ensures that your custom subject cannot produce data faster than the engine can consume it, preventing unbounded memory growth in user-defined connectors.

### How does backpressure interact with autocommit_duration_ms?

These parameters work together to create a two-dimensional memory bound. The `autocommit_duration_ms` parameter groups incoming data into time-based mini-batches, determining the granularity of processing units. The `max_backlog_size` parameter limits how many rows from these mini-batches can exist in memory simultaneously. By tuning both, you control both the temporal window of data visibility and the absolute row count, ensuring that memory usage remains bounded by the product of your batch size and backlog limit rather than the total stream volume.