How to Implement Custom Connectors in Pathway for Unsupported Data Sources
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 (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:
- User code instantiates a concrete subclass of
ConnectorSubject pw.io.python.readwraps the subject in aGenericDataSourceand spawns a dedicated thread viaConnectorSubject.start- Inside
run(), the implementation pulls data from the external system and callsself.next(...)(or deprecated variantsnext_json,next_str,next_bytes) - Each call pushes data into an internal
Queueas(PythonConnectorEventType.INSERT, key, values) - 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.
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. You must implement the run() method to open the external source and repeatedly push records using self.next().
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.
# 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. The following annotated version demonstrates production patterns:
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– ContainsConnectorSubject(lines 49-88), thereadfactory (lines 108-119), and_create_python_datasourcehelper (lines 31-74) that constructs the internalDataStorageandDataFormatobjectsexamples/projects/custom-python-connector-twitter/twitter_connector_example.py– Runnable reference implementation showing real-world HTTP streaming integrationintegration_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.ConnectorSubjectand implement therun()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 typedTablewith a specified schema - Handle cleanup in
on_stop()to prevent resource leaks when pipelines terminate - Control performance via
autocommit_duration_msandmax_backlog_sizeparameters - Reference the Twitter example at
examples/projects/custom-python-connector-twitter/twitter_connector_example.pyfor 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 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 which validates the generic connector interface against mock HTTP endpoints.
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 →