# How to Implement Custom Data Source Connectors for Pathway Pipelines

> Learn to implement custom data source connectors for Pathway pipelines by subclassing BaseConnector or CustomConnector. Easily integrate your data into streaming applications with pw.io.custom.

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

---

**You can implement custom data source connectors for Pathway pipelines by subclassing `pw.io.base.BaseConnector` or `pw.io.custom.CustomConnector`, implementing a `read()` method that returns a `pw.Table`, and consuming it via `pw.io.custom.read()` or `pw.io.custom.stream()` in your pipeline.**

Pathway pipelines process data through **connectors**—thin wrappers that transform arbitrary data sources into `pw.Table` objects that the Pathway engine can process. While the `pathwaycom/llm-app` repository provides many built-in connectors like `pw.io.fs.read` and `pw.io.http.rest_connector`, you often need to implement custom data source connectors for Pathway pipelines when integrating proprietary APIs, internal databases, or specialized message queues.

## Understanding the Pathway Connector Architecture

The connector architecture in Pathway follows a simple contract. A connector class must implement a `read()` method that returns a `pw.Table` containing the data from your source. The Pathway engine then treats this table as a node in the computation graph, automatically tracking changes and propagating updates downstream.

You can base your implementation on `pw.io.base.BaseConnector` for full control, or use `pw.io.custom.CustomConnector` for simplified streaming scenarios. The key distinction is that `read()` should return a static snapshot for batch processing, while continuous sources should use `pw.io.custom.stream()` with an iterator that yields rows indefinitely.

## Step-by-Step Implementation Guide

### Step 1: Create a Connector Class

Subclass `BaseConnector` and initialize it with configuration parameters specific to your data source. Store connection strings, authentication tokens, or client instances as instance attributes.

```python
import pathway as pw
from pathway.io.base import BaseConnector

class DatabaseConnector(BaseConnector):
    def __init__(self, *, connection_string: str, table_name: str):
        self.connection_string = connection_string
        self.table_name = table_name
        self.client = self._create_client()

```

### Step 2: Implement the Read Method

The `read()` method must fetch data from your source and convert it into a Pathway table. Use `pw.from_iterable()` to create a table from Python iterables, or construct rows programmatically.

```python
    def read(self) -> pw.Table:
        # Fetch raw records from your source

        records = self.client.query(f"SELECT * FROM {self.table_name}")
        
        # Convert to list of dictionaries matching your schema

        data = [{"id": r.id, "content": r.content} for r in records]
        
        return pw.from_iterable(data)

```

### Step 3: Expose the Connector Factory

Create a convenience function that instantiates your connector and returns the table, mimicking the signature of built-in Pathway IO functions like `pw.io.fs.read`.

```python
def read_database_table(connection_string: str, table_name: str) -> pw.Table:
    """Read a SQL table as a Pathway table."""
    connector = DatabaseConnector(
        connection_string=connection_string,
        table_name=table_name
    )
    return pw.io.custom.read(connector)

```

### Step 4: Integrate into Your Pipeline

Use your custom connector exactly like built-in sources. Chain Pathway operators such as `.select()`, `.flatten()`, or vector embeddings, then execute with `pw.run()`.

```python
import pathway as pw
from my_connectors import read_database_table

# Create table from custom source

documents = read_database_table(
    connection_string="postgresql://localhost/db",
    table_name="documents"
)

# Process with Pathway operators

processed = documents.select(
    text=documents.content,
    length=pw.apply(len, documents.content)
)

# Execute the pipeline

pw.run()

```

## Code Examples for Custom Data Source Connectors

### Batch Connector for REST APIs

This example implements a connector that fetches JSON data from a REST API endpoint, handling authentication and converting the response into a typed Pathway table.

```python

# my_connectors.py

import requests
import pathway as pw
from pathway.io.base import BaseConnector


class JsonApiConnector(BaseConnector):
    """Fetches a list of JSON objects from a REST endpoint."""
    def __init__(self, *, url: str, auth_token: str | None = None):
        self.url = url
        self.headers = {"Authorization": f"Bearer {auth_token}"} if auth_token else {}

    def read(self) -> pw.Table:
        response = requests.get(self.url, headers=self.headers, timeout=10)
        response.raise_for_status()
        data = response.json()                # → list[dict]

        return pw.from_iterable(data)         # Pathway Table


# usage in a pipeline (e.g., in app.py)

from my_connectors import JsonApiConnector
import pathway as pw

def run(...):
    # 1️⃣ Instantiate the connector with runtime parameters

    connector = JsonApiConnector(url="https://api.example.com/items", auth_token=os.getenv("API_KEY"))

    # 2️⃣ Pull the data as a Pathway table

    raw_table = pw.io.custom.read(connector)

    # 3️⃣ Optional schema enforcement

    class ItemSchema(pw.Schema):
        id: int
        name: str
        price: float

    typed = raw_table.with_columns(**ItemSchema.typehints())
    
    # 4️⃣ Continue with usual Pathway operators

    enriched = typed.select(
        price_usd = pw.apply(lambda p: p * 1.0, typed.price)
    )
    pw.run()

```

### Streaming Connector for Kafka

For continuous data sources like message queues, implement an iterator that yields rows indefinitely and use `pw.io.custom.stream()` instead of `read()`.

```python

# my_kafka_connector.py

import json
from confluent_kafka import Consumer, KafkaException
import pathway as pw
from pathway.io.base import BaseConnector


class KafkaConnector(BaseConnector):
    """Continuously streams JSON messages from a Kafka topic."""
    def __init__(self, *, bootstrap_servers: str, group_id: str, topic: str):
        self.consumer = Consumer({
            "bootstrap.servers": bootstrap_servers,
            "group.id": group_id,
            "auto.offset.reset": "earliest",
        })
        self.consumer.subscribe([topic])

    def read(self) -> pw.Table:
        # `pw.io.custom.stream` expects an iterator that yields rows.

        def iter_messages():
            while True:
                msg = self.consumer.poll(1.0)
                if msg is None:
                    continue
                if msg.error():
                    raise KafkaException(msg.error())
                yield json.loads(msg.value())
        return pw.from_iterable(iter_messages())


# usage

from my_kafka_connector import KafkaConnector
import pathway as pw

def run():
    conn = KafkaConnector(
        bootstrap_servers="localhost:9092",
        group_id="demo-consumer",
        topic="documents",
    )
    # Continuous streaming table

    kafka_table = pw.io.custom.stream(conn)

    # Example: split text, embed, and index

    from pathway.xpacks.llm.embedders import OpenAIEmbedder
    embedder = OpenAIEmbedder(api_key=os.getenv("OPENAI_API_KEY"))
    indexed = kafka_table + kafka_table.select(embedding=embedder(kafka_table.text))
    pw.run()

```

### Factory Wrapper for Simplified Usage

To match the ergonomic API of built-in connectors like `pw.io.fs.read`, expose a thin factory function:

```python

# my_connectors.py (continued)

def json_api(url: str, token: str | None = None) -> pw.Table:
    """Convenient wrapper that mimics pw.io.http.rest_connector."""
    return pw.io.custom.read(JsonApiConnector(url=url, auth_token=token))

```

Now the pipeline can use your connector exactly like native Pathway IO:

```python
raw = json_api(url="https://api.example.com/items", token=os.getenv("API_KEY"))

```

## Key Implementation Patterns from the Pathway LLM App

The `pathwaycom/llm-app` repository demonstrates how connectors integrate into real-world pipelines. In [`templates/unstructured_to_sql_on_the_fly/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/unstructured_to_sql_on_the_fly/app.py) (lines 221-226), the built-in HTTP REST connector is instantiated and wired into the graph:

```python
query, response_writer = pw.io.http.rest_connector(
    host=host,
    port=port,
    schema=NLQuerySchema,
    autocommit_duration_ms=50,
    delete_completed_queries=True,
)

```

Similarly, [`templates/drive_alert/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/drive_alert/app.py) (lines 163-167) shows a source connector pattern using Google Drive:

```python
files = pw.io.gdrive.read(
    object_id=object_id,
    service_user_credentials_file=service_user_credentials_file,
    refresh_interval=30,
)

```

Your custom connector follows the same pattern: instantiate the connector class, obtain a `pw.Table` via `pw.io.custom.read()` or `pw.io.custom.stream()`, then chain Pathway operators like `.select()`, `pw.apply()`, or vector embeddings before calling `pw.run()`.

## Summary

- **Subclass `BaseConnector`** from `pathway.io.base` (or use `CustomConnector`) and implement the `read()` method to return a `pw.Table`.
- **Use `pw.io.custom.read()`** for batch sources and **`pw.io.custom.stream()`** for continuous streaming sources like message queues.
- **Convert raw data** using `pw.from_iterable()` to transform Python iterables into Pathway tables with proper schema enforcement via `pw.Schema` subclasses.
- **Follow repository patterns** demonstrated in [`templates/unstructured_to_sql_on_the_fly/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/unstructured_to_sql_on_the_fly/app.py) and [`templates/drive_alert/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/drive_alert/app.py) to integrate your connector into the Pathway computation graph.
- **Execute** the pipeline with `pw.run()` to enable Pathway's incremental processing and automatic change propagation on your custom data source.

## Frequently Asked Questions

### What is the difference between `pw.io.custom.read` and `pw.io.custom.stream`?

**`pw.io.custom.read`** is designed for batch sources where the connector fetches a finite dataset and returns a static table, similar to `pw.io.fs.read`. **`pw.io.custom.stream`** is for continuous sources like Kafka or message queues where the connector yields an infinite iterator; Pathway treats this as a live stream and automatically updates downstream tables when new rows arrive. Choose `read` for one-time snapshots and `stream` for real-time event processing.

### How do I enforce a schema on data from my custom connector?

After obtaining a table from your custom connector using `pw.io.custom.read()`, you can enforce types by defining a `pw.Schema` subclass and applying it with `with_columns()`. For example:

```python
class ItemSchema(pw.Schema):
    id: int
    name: str

typed_table = raw_table.with_columns(**ItemSchema.typehints())

```

This ensures that downstream operators receive correctly typed columns, matching the pattern used in [`templates/unstructured_to_sql_on_the_fly/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/unstructured_to_sql_on_the_fly/app.py) where `NLQuerySchema` validates HTTP REST inputs.

### Can I use retry strategies and caching with custom connectors?

Yes. While the connector handles data ingestion, you can wrap external API calls within the connector's `read()` method using Pathway's UDF utilities like `pw.udfs.ExponentialBackoffRetryStrategy` or `pw.udfs.DefaultCache`. For example, when calling an external API inside your connector, instantiate the retry strategy and apply it to the fetch operation. This mirrors the error handling patterns seen in the LLM app templates where `OpenAIChat` models use `retry_strategy` and `cache_strategy` parameters.

### Where should I place my custom connector code in the project structure?

Place your custom connector in a separate module within your project (e.g., [`my_connectors.py`](https://github.com/pathwaycom/llm-app/blob/main/my_connectors.py) or [`connectors/kafka_custom.py`](https://github.com/pathwaycom/llm-app/blob/main/connectors/kafka_custom.py)), not inside the Pathway library itself. Import this module in your main [`app.py`](https://github.com/pathwaycom/llm-app/blob/main/app.py) or pipeline entry point, instantiate the connector, and pass it to `pw.io.custom.read()` or `pw.io.custom.stream()`. This approach follows the architecture of the `pathwaycom/llm-app` templates where connectors are instantiated in the main application file (e.g., [`templates/drive_alert/app.py`](https://github.com/pathwaycom/llm-app/blob/main/templates/drive_alert/app.py) lines 163-167) and then chained with downstream processing steps.