# How to Migrate Existing Batch ETL Processes to a Streaming Architecture Using Pathway

> Migrate batch ETL to streaming architecture with Pathway. Effortlessly switch connector mode retain transformation logic for real-time data processing.

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

---

**You can migrate existing batch ETL processes to a streaming architecture using Pathway by changing the connector mode to `"streaming"` and adding an autocommit duration, while keeping your transformation logic unchanged.**

Pathway enables seamless migration from batch to streaming ETL through a unified Python API that works identically in both execution modes. The framework's Rust-based engine handles the complexity of incremental computation, allowing you to convert static data pipelines into real-time streaming architectures with minimal code changes.

## Switching Connectors from Batch to Streaming Mode

The primary migration step involves reconfiguring your input connectors to watch for live data rather than processing static files once.

### S3 and MinIO Connectors

For S3-compatible storage, set `mode="streaming"` when creating the source table. This instructs Pathway to monitor the bucket continuously for new objects or changes.

In [`integration_tests/s3/test_s3_streaming.py`](https://github.com/pathwaycom/pathway/blob/main/integration_tests/s3/test_s3_streaming.py), the test suite demonstrates this pattern:

```python

# integration_tests/s3/test_s3_streaming.py

table = create_table_for_storage(
    storage_type, s3_path, "plaintext", mode="streaming"
)  # ← streaming mode enabled

```

The same approach applies to MinIO and other S3-compatible object stores.

### Delta Lake Connectors

For Delta Lake tables, use `autocommit_duration_ms` to enable streaming-like behavior. This parameter tells the connector to periodically check for new commits and emit updates incrementally.

As shown in [`examples/projects/kafka-alternatives/minio-ETL/etl.py`](https://github.com/pathwaycom/pathway/blob/main/examples/projects/kafka-alternatives/minio-ETL/etl.py):

```python

# examples/projects/kafka-alternatives/minio-ETL/etl.py

timestamps_timezone_1 = pw.io.deltalake.read(
    base_path + "timezone1",
    schema=InputStreamSchema,
    s3_connection_settings=s3_connection_settings,
    autocommit_duration_ms=100,   # <- streaming-like commits every 100ms

)

```

## Reusing Existing Transformation Logic

Pathway's incremental computation engine ensures that most batch transformations work unchanged in streaming mode.

### Stateless Transformations

Operations like `select`, `filter`, and `map` run per-row and require no modifications. The engine automatically applies these transformations to new records as they arrive.

### Stateful Transformations

Joins, windowed aggregations, and reducers such as `pw.reducers.sum()` automatically maintain state across updates. The same aggregation logic used in batch processing continues to compute correct results incrementally without code changes.

## Configuring Streaming Sinks

Output connectors like `pw.io.jsonlines.write`, `pw.io.deltalake.write`, or `pw.io.csv.write` function identically in streaming mode. The sink receives each new batch of rows as they become available and writes them incrementally.

## Running the Pipeline

After defining the graph, call `pw.run()` to start continuous processing. In streaming mode, this call blocks and processes data indefinitely until the process is terminated.

```python
import pathway as pw

# Define your streaming pipeline here

# ...

pw.run()  # Blocks and processes continuously

```

## Complete Migration Examples

### Example 1: Converting a Batch CSV Reader to Streaming

```python
import pathlib
import pathway as pw

# 1️⃣ Batch version (original)

batch_table = pw.io.csv.read(
    pathlib.Path("data/input/"),
    schema=pw.Schema(value=int),
)

# 2️⃣ Streaming version – add mode="streaming"

stream_table = pw.io.csv.read(
    pathlib.Path("data/input/"),
    schema=pw.Schema(value=int),
    mode="streaming",          # <-- enable streaming

)

# 3️⃣ Transform (identical for both)

filtered = stream_table.filter(stream_table.value > 0)
result   = filtered.reduce(total=pw.reducers.sum(stream_table.value))

# 4️⃣ Write out continuously

pw.io.jsonlines.write(result, "output/result.jsonl")

# 5️⃣ Start the engine

pw.run()

```

### Example 2: Streaming Delta Lake ETL with MinIO

```python
import pathway as pw

class Event(pw.Schema):
    ts: str
    payload: str

# Read from Delta Lake in streaming mode (autocommit every 200ms)

events = pw.io.deltalake.read(
    "s3://my-bucket/events",
    schema=Event,
    s3_connection_settings={"endpoint_url": "https://minio.local"},
    autocommit_duration_ms=200,
)

# Simple enrichment – add a parsed timestamp

enriched = events.select(
    ts = pw.this.ts.dt.strptime("%Y-%m-%d %H:%M:%S"),
    payload = pw.this.payload,
)

# Write back to Delta Lake as an incremental table

pw.io.deltalake.write(
    enriched,
    "s3://my-bucket/enriched",
    s3_connection_settings={"endpoint_url": "https://minio.local"},
    min_commit_frequency=200,
)

pw.run()

```

### Example 3: S3 Streaming with JSONLines Output

```python
import pathlib
import pathway as pw
from pathway.tests.utils import ExceptionAwareThread, FileLinesNumberChecker, wait_result_with_checker

# Streaming source watching an S3 bucket

source = pw.io.s3.read("s3://bucket/input/", mode="streaming", format="text")

# Write each new line as a JSON object

output_path = pathlib.Path("tmp/output.jsonl")
pw.io.jsonlines.write(source, output_path)

# Optional: monitor the file size in a separate thread (as the test does)

def monitor():
    wait_result_with_checker(FileLinesNumberChecker(output_path, expected_lines=10), timeout=60)

t = ExceptionAwareThread(target=monitor)
t.start()
pw.run()          # blocks, processing S3 changes forever

t.join()

```

## Summary

- **Change the connector mode**: Add `mode="streaming"` to S3/MinIO connectors or `autocommit_duration_ms` to Delta Lake sources to enable continuous data ingestion.
- **Reuse transformation logic**: Stateless operations (`filter`, `select`, `map`) and stateful aggregations (`reduce`, `join`) work identically in batch and streaming modes.
- **Configure streaming sinks**: Use `pw.io.jsonlines.write`, `pw.io.deltalake.write`, or similar connectors to output incremental results.
- **Run continuously**: Execute `pw.run()` to start the Rust-based engine, which handles incremental computation and blocks until termination.

## Frequently Asked Questions

### Do I need to rewrite my transformation logic when migrating from batch to streaming?

No. Pathway uses the same Python API for both batch and streaming execution. Transformations like `filter`, `select`, `join`, and `reduce` automatically adapt to incremental updates. The engine recomputes only affected rows when new data arrives, so your existing logic remains valid.

### How does Pathway handle late-arriving data in streaming ETL pipelines?

Pathway's Rust engine tracks data dependencies and maintains state automatically for stateful operations like joins and windowed aggregations. When late records arrive, the engine updates the relevant downstream computations incrementally. You do not need to implement custom watermarking logic; the framework handles consistency and correctness guarantees internally.

### What is the difference between `mode="streaming"` and `autocommit_duration_ms`?

`mode="streaming"` is used with connectors like S3 and MinIO to continuously monitor for new files or changes. `autocommit_duration_ms` is specific to Delta Lake connectors and controls how frequently the connector checks for new commits and emits updates. Both parameters enable streaming behavior but apply to different storage backends.

### Can I test streaming pipelines locally before deploying to production?

Yes. Pathway supports local testing using the same `pw.run()` execution model. You can use local file connectors with `mode="streaming"`, or use the `pw.io.python.read` connector to feed test data programmatically. The test suite in [`python/pathway/tests/utils.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/tests/utils.py) provides utilities like `assert_stream_equality` to validate incremental behavior during development.