# How to Optimize Pathway Performance for High-Throughput Real-Time Streaming Data

> Optimize Pathway performance for high-throughput real-time streaming data. Learn to tune connectors, scale horizontally, use Rust UDFs, and manage memory for peak efficiency.

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

---

**To optimize Pathway for high-throughput real-time streaming, configure connectors in streaming mode with tuned `max_backlog_size` and `autocommit_duration_ms`, scale horizontally using `pathway spawn`, and replace Python UDFs with built-in Rust transformations while applying temporal cutoffs to bound memory usage.**

Pathway is a high-performance data processing framework that combines Python ergonomics with a Rust execution engine to handle demanding real-time workloads. In the `pathwaycom/pathway` repository, every component—from the Kafka connectors to the windowing operators—is designed to maximize throughput while maintaining sub-millisecond latency. This guide explains how to optimize Pathway performance for high-throughput real-time streaming data by leveraging its architecture and tuning key configuration parameters.

## Architectural Foundations That Enable Speed

Pathway’s performance stems from a hybrid architecture where **hot paths execute in compiled Rust** while Python provides the API layer. Understanding these foundations is essential before tuning specific parameters.

### The Rust Engine Core

All core data-flow operators run in compiled Rust, not Python. This eliminates the Global Interpreter Lock (GIL) bottleneck and provides deterministic, low-latency execution. The bridge between Python and Rust appears in [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py), where callbacks from the Rust core trigger Python-side updates. This architecture guarantees **high-throughput** and **sub-millisecond** per-record processing.

### Streaming-Mode Connectors

Every I/O connector in Pathway supports a **streaming mode** that continuously polls for new data rather than batch-loading static files. In [`python/pathway/io/kafka/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/kafka/__init__.py), the connector signature exposes `mode: Literal["streaming","static"] = "streaming"`, ensuring the pipeline remains alive and avoids batch-load latency spikes.

### Back-Pressure via max_backlog_size

The engine implements **back-pressure** to pause source reads when internal buffers fill, preventing out-of-memory (OOM) errors during traffic bursts. The `max_backlog_size` parameter appears in every connector implementation—for example, `max_backlog_size: int | None = None` in [`python/pathway/io/csv/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/csv/__init__.py)—allowing you to set a realistic ceiling (e.g., 10,000 rows) that stabilizes latency under load.

### Multi-Process Worker Scaling

Pathway spawns multiple OS processes, each with its own Rust engine thread pool, to achieve linear scaling on multi-core machines. The CLI logic in [`python/pathway/cli.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/cli.py) parses the `--processes` argument and computes total workers, enabling distribution of the dataflow graph across CPU cores.

### Built-in vs Python Transformations

Operations like `window_by`, `asof_now_join`, and reducers such as `pw.reducers.sum()` execute entirely in Rust. These **built-in transformations** are 10–100× faster than Python UDFs and allow the scheduler to fuse operators, reducing data movement overhead across the Python-Rust boundary.

## Practical Tuning Guidelines

Optimizing throughput requires balancing latency, memory, and CPU utilization through specific parameter configurations.

### Connector Configuration for Streaming Sources

Always set `mode="streaming"` for real-time feeds (Kafka, Redpanda, S3, filesystem). Tune `autocommit_duration_ms` to 200–500 ms for low-latency commits—avoiding excessive commit overhead—while setting `max_backlog_size` to prevent memory spikes during bursty production. These parameters are documented in the connector implementations at [`python/pathway/io/kafka/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/kafka/__init__.py) and [`python/pathway/io/csv/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/csv/__init__.py).

### Horizontal Scaling with CLI

For CPU-bound pipelines, launch multiple processes using the spawn command:

```bash
pathway spawn -n 8 python my_pipeline.py

```

The CLI logic in [`python/pathway/cli.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/cli.py) (lines 245–251) handles worker distribution automatically. Start with **2 × CPU-core count** and monitor utilization to find the optimal worker density.

### Memory Management Strategies

Reduce memory pressure through three techniques:

- **Use `asof_now_join`** when joining a high-throughput stream against a slowly-changing dimension table. This stores only the left side of the join in memory, dramatically reducing the state footprint.
- **Apply temporal cutoffs** via `behavior=pw.temporal.common_behavior(cutoff=60)` to discard window state after 60 seconds, bounding memory usage for sliding and tumbling windows.
- **Mark deterministic UDFs** with `@pw.udf(deterministic=True)` to allow the engine to skip redundant computations, saving both memory and CPU cycles.

## Code Implementation Examples

### Configuring High-Throughput Kafka Sources

The following configuration minimizes parsing overhead while enforcing back-pressure:

```python
import pathway as pw

rdkafka_settings = {
    "bootstrap.servers": "localhost:9092",
    "security.protocol": "plaintext",
}

events = pw.io.kafka.read(
    rdkafka_settings,
    topic="events",
    format="raw",                    # Minimal parsing overhead

    mode="streaming",                # Keep pipeline alive

    autocommit_duration_ms=200,      # Fast commits for low latency

    max_backlog_size=10_000,         # Back-pressure on bursts

    name="kafka_events",
)

# Built-in Rust transformation: count events per second

counts = (
    events.window_by(pw.temporal.tumbling_window(duration=1))
          .reduce(pw.reducers.count())
)

```

The parameters `mode`, `autocommit_duration_ms`, and `max_backlog_size` are defined in [`python/pathway/io/kafka/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/kafka/__init__.py).

### Scaling to Multiple Processes

To utilize all available cores, wrap your pipeline logic in a script and spawn workers:

```bash

# Launch with 8 OS processes (each with independent Rust engine threads)

pathway spawn -n 8 python high_throughput_pipeline.py

```

The spawn logic in [`python/pathway/cli.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/cli.py) automatically partitions the dataflow graph across workers without requiring code changes.

### Optimizing Joins and Aggregations

Replace standard joins with `asof_now_join` for look-up scenarios, and bound window state with cutoffs:

```python

# Stream of device telemetry

readings = pw.io.kafka.read(..., schema=DeviceSchema)

# Small static metadata table

metadata = pw.debug.table_from_rows(
    schema=MetadataSchema,
    rows=[...],
    mode="static"
)

# Memory-efficient join: stores only readings side

enriched = readings.asof_now_join(
    metadata, 
    readings.device_id == metadata.device_id
)

# Bounded aggregation: discard state after 60 seconds

windowed = (
    readings.window_by(
        window=pw.temporal.sliding(duration=10, hop=2),
        behavior=pw.temporal.common_behavior(cutoff=60)
    )
    .reduce(pw.reducers.mean(readings.temperature))
)

```

The `asof_now_join` implementation and temporal behavior definitions reside in the Pathway standard library, documented in the best-practices guides and implemented in [`python/pathway/stdlib/temporal/temporal_behavior.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/stdlib/temporal/temporal_behavior.py).

## Key Source Files for Performance Tuning

| File | Performance Relevance |
|------|----------------------|
| [`python/pathway/io/kafka/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/kafka/__init__.py) | Streaming mode flags, `max_backlog_size` handling, and commit frequency controls for high-throughput ingestion. |
| [`python/pathway/io/csv/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/csv/__init__.py) | Reference implementation of back-pressure parameters applicable to file-based sources. |
| [`python/pathway/cli.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/cli.py) | Worker scaling logic and `--processes` argument parsing for multi-process deployment. |
| [`python/pathway/internals/datasource.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/datasource.py) | Low-level abstraction propagating `max_backlog_size` to the Rust runtime. |
| [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py) | Python-Rust bridge demonstrating how the engine executes without GIL contention. |

## Summary

- **Enable streaming mode** and tune `max_backlog_size` plus `autocommit_duration_ms` to balance latency and memory safety.
- **Scale horizontally** using `pathway spawn -n <workers>` to distribute load across CPU cores via the Rust engine.
- **Prefer built-in transformations** over Python UDFs for hot paths, and use `asof_now_join` to minimize join state.
- **Apply temporal cutoffs** to window operations to prevent unbounded memory growth in long-running pipelines.

## Frequently Asked Questions

### What is the optimal autocommit_duration_ms for low-latency streaming?

Set `autocommit_duration_ms` between 200 and 500 milliseconds for real-time pipelines. Values below 200 ms increase commit overhead without meaningful latency gains, while values above 1000 ms may delay visibility of processed records. This parameter is available in all streaming connectors, including the Kafka implementation in [`python/pathway/io/kafka/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/kafka/__init__.py).

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

The `max_backlog_size` parameter places a hard ceiling on the number of unprocessed records buffered between a source and the engine. When the buffer reaches this limit, the connector pauses reads, applying **back-pressure** to upstream producers. This prevents memory spikes during traffic bursts and is implemented in [`python/pathway/internals/datasource.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/datasource.py).

### When should I use asof_now_join instead of regular joins?

Use `asof_now_join` when joining a high-throughput stream against a relatively static or slowly-changing table (such as device metadata or user profiles). Unlike standard interval joins that store both sides in state, `asof_now_join` retains only the streaming side, reducing memory usage by orders of magnitude for high-cardinality look-ups.

### How many processes should I spawn for CPU-bound pipelines?

Start with **2 × CPU-core count** and monitor CPU utilization via system metrics. If utilization remains below 80%, increase the count; if overhead from coordination appears, reduce it. The `pathway spawn` command in [`python/pathway/cli.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/cli.py) handles the distribution automatically, allowing linear scaling until I/O or network limits are reached.