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
BaseConnectorfrompathway.io.base(or useCustomConnector) and implement theread()method to return apw.Table. - Use
pw.io.custom.read()for batch sources andpw.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 viapw.Schemasubclasses. - Follow repository patterns demonstrated in
templates/unstructured_to_sql_on_the_fly/app.pyandtemplates/drive_alert/app.pyto 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →