How to Implement Custom Data Source Connectors for Pathway Pipelines

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.

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.

    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.

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().

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.


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


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


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

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 (lines 221-226), the built-in HTTP REST connector is instantiated and wired into the graph:

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 (lines 163-167) shows a source connector pattern using Google Drive:

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

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 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 or connectors/kafka_custom.py), not inside the Pathway library itself. Import this module in your 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 lines 163-167) and then chained with downstream processing steps.

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 →