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

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, the test suite demonstrates this pattern:


# 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:


# 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.

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

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

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

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 provides utilities like assert_stream_equality to validate incremental behavior during development.

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 →