# How to Implement Custom Connectors in Pathway for Unsupported Data Sources

> Learn to implement custom connectors in Pathway for unsupported data sources. Subclass ConnectorSubject, use self.next(), and wrap with read() for streaming data.

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

---

**Implement custom connectors in Pathway by subclassing `pw.io.python.ConnectorSubject`, implementing the `run()` method to push data via `self.next()`, and wrapping it with `pw.io.python.read()` to create a streaming table.**

Pathway provides a Python-based connector API that enables ingestion from any external system not covered by built-in connectors like Kafka, PostgreSQL, or Airbyte. The pathwaycom/pathway repository exposes this extensibility through the abstract `ConnectorSubject` class and the `pw.io.python.read` factory function, which bridges Python I/O code with Pathway's Rust query engine. This architecture isolates blocking I/O operations in Python threads while maintaining high-performance stream processing in the core engine.

## Architecture of Pathway's Python Connector API

The connector framework centers on three core components that handle the boundary between external data sources and Pathway's internal representation.

**ConnectorSubject** is the abstract base class defined in [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py) (lines 49-88). It establishes the contract for custom sources by providing an internal buffer, a dedicated thread, and the abstract `run()` method that must feed data into the buffer via `self.next(...)`. Each call to `next()` enqueues a tuple containing the event type, key, and values that the Rust engine consumes as change events.

**pw.io.python.read** acts as the factory function (lines 108-119 in the same file) that transforms a running subject into a streaming `Table`. This function wires the subject's `start`, `seek`, `read`, and `end` callbacks into a `GenericDataSource` and configures the `DataStorage` and `DataFormat` objects required by the core engine through the internal `_create_python_datasource` helper.

**Data Flow Process**:

1. User code instantiates a concrete subclass of `ConnectorSubject`
2. `pw.io.python.read` wraps the subject in a `GenericDataSource` and spawns a dedicated thread via `ConnectorSubject.start`
3. Inside `run()`, the implementation pulls data from the external system and calls `self.next(...)` (or deprecated variants `next_json`, `next_str`, `next_bytes`)
4. Each call pushes data into an internal `Queue` as `(PythonConnectorEventType.INSERT, key, values)`
5. The Rust core pulls from this queue, materializes rows according to the supplied schema, and makes them available downstream in the Pathway graph

## Step-by-Step Guide to Build a Custom Connector

### Define Your Schema

First, declare a schema that matches the structure of data your connector will emit. Use `pw.Schema` with column definitions to specify primary keys and data types.

```python
import pathway as pw

class TweetSchema(pw.Schema):
    id: int = pw.column_definition(primary_key=True)
    text: str
    created_at: str

```

### Subclass ConnectorSubject

Implement a concrete subclass of `pw.io.python.ConnectorSubject` located at [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py). You must implement the `run()` method to open the external source and repeatedly push records using `self.next()`.

```python
class TwitterSubject(pw.io.python.ConnectorSubject):
    """Ingests live tweets from the Twitter streaming API."""
    
    def __init__(self, bearer_token: str) -> None:
        super().__init__()
        self._client = tweepy.StreamingClient(bearer_token)
        self._client.on_response = self._handle_tweet
    
    def run(self) -> None:
        """Starts the external connection and blocks until data arrives."""
        self._client.sample()  # Starts the streaming connection

    
    def _handle_tweet(self, response) -> None:
        """Callback invoked by the external client for each record."""
        self.next(id=response.data.id, text=response.data.text)
    
    def on_stop(self) -> None:
        """Cleanup resources when the pipeline shuts down."""
        self._client.disconnect()

```

Key implementation requirements:
- **run()**: Must contain the blocking I/O loop that fetches data from your source
- **self.next()**: Accepts keyword arguments matching your schema columns to emit rows
- **on_stop()**: Optional but recommended for graceful shutdown (closing sockets, HTTP sessions, or database connections)

### Wire Everything Together

Use `pw.io.python.read()` to convert your subject into a live table, then consume it like any other Pathway table.

```python

# Instantiate the custom subject

subject = TwitterSubject(bearer_token=os.environ["TWITTER_API_TOKEN"])

# Create the streaming table

tweets = pw.io.python.read(
    subject,
    schema=TweetSchema,
    autocommit_duration_ms=1000,  # Batch commit interval in milliseconds

    max_backlog_size=10_000       # Memory back-pressure limit

)

# Consume the stream

pw.io.csv.write(tweets, "output.csv")
pw.run()

```

## Full Working Example: Streaming Twitter Data

The pathwaycom/pathway repository includes a complete implementation at [`examples/projects/custom-python-connector-twitter/twitter_connector_example.py`](https://github.com/pathwaycom/pathway/blob/main/examples/projects/custom-python-connector-twitter/twitter_connector_example.py). The following annotated version demonstrates production patterns:

```python
import os
import tweepy
import pathway as pw

# ----------------------------------------------------------------------

# 1. Define the schema that matches downstream expectations

# ----------------------------------------------------------------------

class InputSchema(pw.Schema):
    key: int = pw.column_definition(primary_key=True)
    text: str

# ----------------------------------------------------------------------

# 2. Implement the connector subject

# ----------------------------------------------------------------------

class TwitterSubject(pw.io.python.ConnectorSubject):
    """Reads live tweets from the Twitter streaming API."""
    
    def __init__(self) -> None:
        super().__init__()
        self._client = TwitterClient(self)
    
    def run(self) -> None:
        # Opens a random sample stream; blocks until connection closes

        self._client.sample()
    
    def on_stop(self) -> None:
        # Graceful shutdown when Pathway finishes

        self._client.disconnect()

class TwitterClient(tweepy.StreamingClient):
    """Thin wrapper that forwards Twitter events to the subject."""
    
    def __init__(self, subject: TwitterSubject) -> None:
        super().__init__(os.environ["TWITTER_API_TOKEN"])
        self._subject = subject
    
    def on_response(self, response) -> None:
        """Invoked by Tweepy for each tweet received."""
        self._subject.next(key=response.data.id, text=response.data.text)

# ----------------------------------------------------------------------

# 3. Build and run the pipeline

# ----------------------------------------------------------------------

if __name__ == "__main__":
    tweets = pw.io.python.read(TwitterSubject(), schema=InputSchema)
    pw.io.csv.write(tweets, "output.csv")
    
    try:
        pw.run()
    except KeyboardInterrupt:
        print("Pipeline stopped by user.")

```

This example shows the complete lifecycle: schema definition, subject implementation with external client integration, and pipeline execution.

## Key Implementation Files and References

Understanding the source structure helps when debugging or extending custom connectors:

- **[`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py)** – Contains `ConnectorSubject` (lines 49-88), the `read` factory (lines 108-119), and `_create_python_datasource` helper (lines 31-74) that constructs the internal `DataStorage` and `DataFormat` objects
- **[`examples/projects/custom-python-connector-twitter/twitter_connector_example.py`](https://github.com/pathwaycom/pathway/blob/main/examples/projects/custom-python-connector-twitter/twitter_connector_example.py)** – Runnable reference implementation showing real-world HTTP streaming integration
- **[`integration_tests/webserver/test_rest_connector.py`](https://github.com/pathwaycom/pathway/blob/main/integration_tests/webserver/test_rest_connector.py)** – Test suite validating the generic connector interface, useful for verifying custom implementations

## Performance Considerations and Best Practices

**Back-pressure Management**: If your external source produces data faster than Pathway processes it, set `max_backlog_size` in `pw.io.python.read()` to cap memory usage. When the buffer reaches this limit, the subject's thread blocks until the engine catches up.

**Deletion Support**: By default, connectors operate in append-only mode. To enable UPSERT semantics (inserts and updates with deletes), override the `_deletions_enabled` property in your subject to return `True`, or use `SessionType.UPSERT` in internal API calls.

**Commit Frequency**: The `autocommit_duration_ms` parameter controls batching. Lower values (e.g., 100ms) reduce latency but increase overhead; higher values (e.g., 5000ms) improve throughput but increase latency.

**Thread Safety**: All calls to `self.next()` must originate from the thread started by `ConnectorSubject.start()` (inside your `run()` method). Do not invoke `next()` from other threads unless you implement explicit synchronization, as the internal `Queue` is not thread-safe for external producers.

## Summary

- **Subclass `pw.io.python.ConnectorSubject`** and implement the `run()` method to fetch data from any external source
- **Push records** using `self.next(key=value, ...)` which enqueues them for the Rust engine
- **Wrap with `pw.io.python.read()`** to convert the subject into a typed `Table` with a specified schema
- **Handle cleanup** in `on_stop()` to prevent resource leaks when pipelines terminate
- **Control performance** via `autocommit_duration_ms` and `max_backlog_size` parameters
- **Reference the Twitter example** at [`examples/projects/custom-python-connector-twitter/twitter_connector_example.py`](https://github.com/pathwaycom/pathway/blob/main/examples/projects/custom-python-connector-twitter/twitter_connector_example.py) for a complete production pattern

## Frequently Asked Questions

### How do I handle back-pressure in custom Pathway connectors?

Set the `max_backlog_size` parameter in `pw.io.python.read()` to limit the internal queue size. When the buffer reaches this limit, calls to `self.next()` will block until the engine consumes records, preventing unbounded memory growth when the external source outpaces processing capacity.

### Can custom connectors support deletion events (UPSERT semantics)?

Yes, but you must explicitly enable deletions by overriding the `_deletions_enabled` property in your `ConnectorSubject` subclass to return `True`. By default, connectors operate in append-only mode. When deletions are enabled, you can emit removal events by specifying `PythonConnectorEventType.DELETE` in advanced internal APIs, though most users achieve UPSERT semantics through key-based updates via the standard `next()` method.

### What is the difference between `ConnectorSubject` and `ConnectorObserver`?

`ConnectorSubject` (defined in [`python/pathway/io/python/__init__.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/io/python/__init__.py) lines 49-88) is designed for **input** connectors—reading data from external sources into Pathway. `ConnectorObserver` and `AsyncConnectorObserver` handle **output** connectors—writing processed results back to external sinks. For source-only implementations, you only need `ConnectorSubject`.

### How do I test a custom connector before production deployment?

Isolate your connector in a standalone script that calls `pw.run()` with a short `autocommit_duration_ms` (e.g., 100ms) and writes to a local CSV or JSONL sink. Verify that `run()` properly initializes, `self.next()` correctly formats data matching your schema, and `on_stop()` releases resources when interrupted. For automated testing, reference patterns in [`integration_tests/webserver/test_rest_connector.py`](https://github.com/pathwaycom/pathway/blob/main/integration_tests/webserver/test_rest_connector.py) which validates the generic connector interface against mock HTTP endpoints.